From f577834599ac13b65225a5d3329d06880948ab7b Mon Sep 17 00:00:00 2001 From: Simon Binder Date: Tue, 1 Sep 2026 14:50:21 +0200 Subject: [PATCH 01/15] Represent instantiated parameters as result set --- .../evaluator/parameter_evaluator.ts | 724 ++++++------------ .../src/sync_plan/evaluator/result_set.ts | 177 +++++ .../sync_plan/evaluator/result_set.test.ts | 179 +++++ 3 files changed, 570 insertions(+), 510 deletions(-) create mode 100644 packages/sync-rules/src/sync_plan/evaluator/result_set.ts create mode 100644 packages/sync-rules/test/src/sync_plan/evaluator/result_set.test.ts diff --git a/packages/sync-rules/src/sync_plan/evaluator/parameter_evaluator.ts b/packages/sync-rules/src/sync_plan/evaluator/parameter_evaluator.ts index 0a43dd7d8..dffe7ca68 100644 --- a/packages/sync-rules/src/sync_plan/evaluator/parameter_evaluator.ts +++ b/packages/sync-rules/src/sync_plan/evaluator/parameter_evaluator.ts @@ -1,6 +1,5 @@ import { ParameterLookupSource, ScopedParameterLookup, UnscopedParameterLookup } from '../../BucketParameterQuerier.js'; import { ParameterIndexLookupCreator } from '../../BucketSource.js'; -import { HashMap, listEquality, StableHasher } from '../../compiler/equality.js'; import { HydrationState } from '../../HydrationState.js'; import { RequestParameters, SqliteParameterValue, SqliteValue } from '../../types.js'; import { isValidParameterValue } from '../../utils.js'; @@ -14,6 +13,7 @@ import { MapSourceVisitor, visitExpr } from '../expression_visitor.js'; import * as plan from '../plan.js'; import { StreamInput } from './bucket_source.js'; import { PreparedParameterIndexLookupCreator } from './parameter_index_lookup_creator.js'; +import { AsyncJoinLookup, ResultSet, ResultSetColumn, ResultSetElement } from './result_set.js'; /** * Finds bucket parameters for a given request or subscription. @@ -63,7 +63,12 @@ export class RequestParameterEvaluators { /** * Pending parameter values, or their cached outputs. */ - readonly parameterValues: PreparedParameterValue[] + readonly parameterValues: PreparedParameterValue[], + + /** + * The materialized result set into which + */ + readonly resultSet: ResultSet ) {} /** @@ -77,33 +82,10 @@ export class RequestParameterEvaluators { * instead of re-evaluating them on every parameter lookup change. */ clone(): RequestParameterEvaluators { - function cloneValue(value: PreparedParameterValue): PreparedParameterValue { - switch (value.type) { - case 'intersection': - return { type: 'intersection', values: value.values.map(cloneValue) }; - case 'request': - case 'lookup': - case 'cached': - return value; - } - } + const copiedStages = this.lookupStages.map((s) => s.map((e) => e.clone())); + const outputValues = this.parameterValues.map((v) => v.clone()); - function cloneLookup(lookup: PreparedExpandingLookup): PreparedExpandingLookup { - switch (lookup.type) { - case 'parameter': - // We need to clone the instantiation array as well. - return { type: 'parameter', lookup: lookup.lookup, instantiation: lookup.instantiation.map(cloneValue) }; - case 'table_valued': - case 'cached': - return lookup; - } - } - - return new RequestParameterEvaluators( - this.stream, - this.lookupStages.map((stage) => stage.map(cloneLookup)), - this.parameterValues.map(cloneValue) - ); + return new RequestParameterEvaluators(this.stream, copiedStages, outputValues, this.resultSet.clone()); } /** @@ -115,13 +97,34 @@ export class RequestParameterEvaluators { * If dynamic lookups are required to resolve parameters, returns `undefined`. */ partiallyInstantiate(input: PartialInstantiationInput): SqliteParameterValue[][] | undefined { - const helper = new PartialInstantiator(input, this); + try { + // At this point, we can resolve table-valued lookups and parameter values based only on request data. + for (const stage of this.lookupStages) { + for (const element of stage) { + if (element instanceof TableValuedExpandingLookup) { + const outputs = element.read(input.request); + element.wasResolved = true; + this.resultSet.multiply(element.resultSetIndex, outputs); + + this.#checkInstantiable(); + } else { + for (const instantiation of element.instantiation) { + if (instantiation instanceof RequestParameterValue) instantiation.resolveWith(input); + } + } + } + } - this.lookupStages.forEach((stage, stageIndex) => { - stage.forEach((_, indexInStage) => helper.expandingLookupSync(stageIndex, indexInStage)); - }); + for (const parameter of this.parameterValues) { + if (parameter instanceof RequestParameterValue) parameter.resolveWith(input); + } - return helper.tryResolveInstantiation(this.parameterValues)?.map(withoutProvenance); + return this.#readParameters(); + } catch (e) { + if (e === uninstantiableException) return []; + + throw e; + } } /** @@ -130,16 +133,103 @@ export class RequestParameterEvaluators { * Because this needs to lookup parameter indexes, it is asynchronous. */ async instantiate(input: InstantiationInput): Promise { - const helper = new FullInstantiator(input, this); + try { + for (const stage of this.lookupStages) { + for (const lookup of stage) { + if (lookup instanceof ParameterIndexExpandingLookup) { + await this.#instantiateLookup(lookup, input); + } + } + } + + return this.#readParameters()!; + } catch (e) { + if (e === uninstantiableException) return []; + + throw e; + } + } - for (let i = 0; i < this.lookupStages.length; i++) { - // Within a stage, we can resolve lookups concurrently. - await Promise.all(this.lookupStages[i].map((_, j) => helper.expandingLookup(i, j))); + #checkInstantiable() { + if (this.resultSet.length === 0) throw uninstantiableException; + } + + #readParameters(): SqliteParameterValue[][] | undefined { + for (const stage of this.lookupStages) { + for (const element of stage) { + if (!element.wasResolved) { + return undefined; + } + } } - // At this point, all lookups have been resolved and we can synchronously evaluate parameters which might depend on - // those lookups. - return helper.resolveInstantiation(this.parameterValues).map(withoutProvenance); + return this.#readValues(this.parameterValues); + } + + #readValues(values: PreparedParameterValue[]): SqliteParameterValue[][] { + const allInstantiations: SqliteParameterValue[][] = []; + + for (const row of this.resultSet.projectUnique(values.filter((v) => v instanceof LookupParameterValue))) { + allInstantiations.push(this.#evaluateAgainstRow(row, values)); + } + + return allInstantiations; + } + + #evaluateAgainstRow(row: SqliteParameterValue[], projection: PreparedParameterValue[]) { + let lookupIndex = 0; + + return projection.map((v) => { + if (v instanceof LookupParameterValue) { + return row[lookupIndex++]; + } else { + if (!v.resolved) throw new Error('Expected request values to be resolved here'); + return v.resolved; + } + }); + } + + async #instantiateLookup(lookup: ParameterIndexExpandingLookup, input: InstantiationInput) { + const scope = input.hydrationState.getParameterIndexLookupScope(lookup.lookup); + const resolvedLookup = lookup.lookup as PreparedParameterIndexLookupCreator; + + await this.resultSet.joinAsync( + lookup.instantiation.filter((v) => v instanceof LookupParameterValue), + lookup.resultSetIndex, + async (inputs) => { + const bucketStorageLookups = new Map(); + for (const input of inputs) { + bucketStorageLookups.set( + ScopedParameterLookup.normalized( + scope, + UnscopedParameterLookup.normalized(this.#evaluateAgainstRow(input.inputs, lookup.instantiation)) + ), + input + ); + } + + const outputs = await input.source.getParameterSets( + [...bucketStorageLookups.keys()], + `Stream ${this.stream.name} evaluating parameter on ${resolvedLookup.sourceTable.tablePattern}` + ); + + for (const { lookup, rows } of outputs) { + const join = bucketStorageLookups.get(lookup)!; + + for (const row of rows) { + const length = Object.entries(row).length; + const asArray: SqliteParameterValue[] = []; + for (let i = 0; i < length; i++) { + asArray.push(row[i.toString()] as SqliteParameterValue); + } + + join.foundRows.push(asArray); + } + } + } + ); + lookup.wasResolved = true; + this.#checkInstantiable(); } /** @@ -160,7 +250,8 @@ export class RequestParameterEvaluators { engine: ScalarExpressionEngine ) { const mappedStages: PreparedExpandingLookup[][] = []; - const lookupToStage = new Map(); + let amountOfLookups = 0; + const lookupToStage = new Map(); function mapParameterValue(value: plan.ParameterValue): PreparedParameterValue { if (value.type == 'request') { @@ -169,17 +260,14 @@ export class RequestParameterEvaluators { const prepared = engine.prepareEvaluator({ filters: [], outputs: [mapper.transform(value.expr)] }); const instantiation = mapper.instantiation; - return { - type: 'request', - read(request) { - return prepared.evaluate(parametersForRequest(request, instantiation))[0][0]; - } - }; + return new RequestParameterValue( + (request) => prepared.evaluate(parametersForRequest(request, instantiation))[0][0] + ); } else if (value.type == 'lookup') { - const stagePosition = lookupToStage.get(value.lookup)!; - return { type: 'lookup', lookup: stagePosition, resultIndex: value.resultIndex }; + const lookup = lookupToStage.get(value.lookup)!; + return new LookupParameterValue(lookup, value.resultIndex); } else { - return { type: 'intersection', values: mapParameterValues(value.values) }; + throw new Error('TODO: intersection'); } } @@ -194,14 +282,14 @@ export class RequestParameterEvaluators { for (const lookup of stage) { const index = mappedStage.length; - lookupToStage.set(lookup, { stage: stageIndex, index }); + let resolved: PreparedExpandingLookup; if (lookup.type == 'parameter') { - mappedStage.push({ - type: 'parameter', - lookup: input.preparedLookups.get(lookup.lookup)!, - instantiation: mapParameterValues(lookup.instantiation) - }); + resolved = new ParameterIndexExpandingLookup( + amountOfLookups++, + input.preparedLookups.get(lookup.lookup)!, + mapParameterValues(lookup.instantiation) + ); } else { // Create an expression like SELECT FROM table_valued() WHERE const mapInputs = mapExternalDataToInstantiation(); @@ -222,430 +310,107 @@ export class RequestParameterEvaluators { filters: lookup.filters.map((e) => visitExpr(mapOutputs, e, null)) }); - mappedStage.push({ - type: 'table_valued', - read(request) { - return [ - ...filterParameterRows(prepared.evaluate(parametersForRequest(request, mapInputs.instantiation))) - ]; - } - }); + resolved = new TableValuedExpandingLookup(amountOfLookups++, (request) => [ + ...filterParameterRows(prepared.evaluate(parametersForRequest(request, mapInputs.instantiation))) + ]); } - } - } - return new RequestParameterEvaluators(stream, mappedStages, mapParameterValues(values)); - } -} - -class PartialInstantiator { - constructor( - protected readonly input: I, - protected readonly evaluators: RequestParameterEvaluators - ) {} - - tryResolveInstantiation(params: PreparedParameterValue[]): ParameterValueWithRow[][] | undefined { - const stages = this.evaluators.lookupStages; - let hasUninstantiatedStage = false; - for (let stageIndex = 0; stageIndex < stages.length; stageIndex++) { - const stage = stages[stageIndex]; - for (let indexInStage = 0; indexInStage < stage.length; indexInStage++) { - const resolvedValues = this.expandingLookupSync(stageIndex, indexInStage); - if (resolvedValues == null) { - // Requires an asynchronous lookup to instantiate. - hasUninstantiatedStage = true; - continue; - } - - if (resolvedValues.length == 0) { - // Empty lookup stages make the entire graph uninstantiable, even if they're not used as a parameter. The - // reason for that is that queries like `WHERE 'static_value' IN (SELECT name FROM users WHERE id = auth.user_id())` - // are implemented as lookup stages, so we can't ignore them. - // Note that there is no construct like `OR` in a querier lookup (those always get compiled into separate - // queries), so any stage being empty guarantees that everything is uninstantiable. - return []; - } + lookupToStage.set(lookup, resolved); } } - if (hasUninstantiatedStage) { - return undefined; - } - - // If we got to this point, all stages have been resolved. So we can resolve parameters without further async work. - return [...this.resolveInputs(params)]; - } - - protected *resolveInputs(params: PreparedParameterValue[]): Generator { - const parameterValues = params.map((_, index) => { - const cached = this.parameterSync(params, index); - if (cached == null) { - // This method is only called for inputs from an earlier stage, which should have been resolved at this point. - throw new Error('Should have been able to resolve parameter from earlier stage synchronously.'); - } - return cached; - }); - - yield* mergeValueCombinations(parameterValues); + const rs = new ResultSet(amountOfLookups); + return new RequestParameterEvaluators(stream, mappedStages, mapParameterValues(values), rs); } +} - /** - * If possible, evaluates an element in an array of parameter values and replaces the parameter with a marker - * indicating it as cached. - */ - parameterSync(parent: PreparedParameterValue[], index: number): ParameterValueWithRow[] | undefined { - const current = parent[index]; - if (current.type === 'cached') { - return current.values; - } else if (current.type === 'intersection') { - const columns: ParameterValueWithRow[][] = []; - for (let i = 0; i < current.values.length; i++) { - const evaluated = this.parameterSync(current.values, i); - if (evaluated == null) { - return undefined; // Can't evaluate sub-parameter - } - columns.push(evaluated); - } - - // For the most part, this just needs to find an intersection of values present in all columns. It gets more - // complicated for rows with provenance, however. For those. we need to ensure we find an intersection of values - // with compatible source rows. For example, consider an intersection of the same parameter lookup with columns - // `c1` and `c2`, and assume that we had the following rows: - // - // 1. Row {c1: 'a', c2: 'a'} - // 2. Row {c1: 'a', c2: 'b'} - // - // The intersection of this has one value: `a`, with a provenance of Row 1. To achieve this, we re-create rows - // by tracking a canonical value per row. If we see another value in the same row, we know that row can't - // contribute to the intersection because it has different values for `c1` and `c2`. - const poison = Symbol('poison'); - const valuesByResultSet = new Map>(); - const completedRows = new Map>(); - function markRowAsCompleted(resultSet: symbol, rowid: number) { - let rowids = completedRows.get(resultSet); - if (rowids == null) { - rowids = new Set(); - completedRows.set(resultSet, rowids); - } - - rowids.add(rowid); - } - - // Eliminate rows with conflicting values. - for (const column of columns) { - nextValue: for (const value of column) { - const row = value.directOrigin; - if (row) { - let forResultSet = valuesByResultSet.get(row.resultSet); - if (forResultSet == null) { - forResultSet = new Map(); - valuesByResultSet.set(row.resultSet, forResultSet); - } - - const existingValue = forResultSet.get(row.row); - if (existingValue != null && existingValue != value.value) { - forResultSet.set(row.row, poison); - markRowAsCompleted(row.resultSet, row.row); - continue nextValue; - } else { - forResultSet.set(row.row, value.value); - } - } - } - } - - function shouldSkipValue(value: ParameterValueWithRow): boolean { - if (value.directOrigin) { - const { resultSet, row } = value.directOrigin; - - const ignoredRowIds = completedRows.get(resultSet); - if (ignoredRowIds?.has(row)) return true; - - const valuesForResultSet = valuesByResultSet.get(resultSet); - if (valuesForResultSet == null) return false; - - if (valuesForResultSet.get(row) != value.value) return true; - } - - return false; - } - - let intersection: Map | null = null; - for (const column of columns) { - if (intersection == null) { - intersection = new Map(); - - for (const value of column) { - if (shouldSkipValue(value)) continue; - - const existing = intersection.get(value.value); - if (existing != null) { - existing.push(value.provenance); - } else { - intersection.set(value.value, [value.provenance]); - } - - for (const { resultSet, row } of value.provenance) { - // Any other value derived from this row must have the same value (otherwise we would have eliminated it). - // So we don't have to consider this row again. - markRowAsCompleted(resultSet, row); - } - } - } else { - const unmatchedValues = new Set(intersection.keys()); - - for (const value of column) { - const existing = intersection.get(value.value); - if (existing == null) { - // Value not in intersection, ignore. - } else { - unmatchedValues.delete(value.value); - - if (!shouldSkipValue(value)) { - // An intersection value is derived from all inputs, so we track them all as provenance. - existing.push(value.provenance); - } - } - } - - for (const unmatched of unmatchedValues) { - // Values in intersection before, but not in evaluated - intersection.delete(unmatched); - } - } +/** + * An internal exception thrown when no instantiation exists for a parameter. + * + * This is an exception to allow aborting the evaluator early. + */ +const uninstantiableException = Symbol.for('uninstantiable'); - if (intersection!.size == 0) { - // Empty intersection, we don't even need to evaluate the rest. - break; - } - } +export type PreparedExpandingLookup = TableValuedExpandingLookup | ParameterIndexExpandingLookup; - let values: ParameterValueWithRow[] = []; - if (intersection) { - intersection.forEach((provenances, value) => { - for (const provenance of provenances) { - values.push({ value, provenance }); - } - }); - } +abstract class BasePreparedExpandingLookup implements ResultSetElement { + wasResolved = false; - parent[index] = { type: 'cached', values }; - return values; - } else if (current.type === 'lookup') { - const resolvedLookup = this.expandingLookupSync(current.lookup.stage, current.lookup.index); - if (resolvedLookup) { - const values = resolvedLookup.map((row) => row[current.resultIndex]); - parent[index] = { type: 'cached', values }; - return values; - } - } else if (current.type === 'request') { - const value = current.read(this.input.request); - const values: ParameterValueWithRow[] = isValidParameterValue(value) - ? [ - { - value, - provenance: [] - } - ] - : []; + constructor(readonly resultSetIndex: number) {} - parent[index] = { type: 'cached', values }; - return values; - } + abstract clone(): BasePreparedExpandingLookup; +} - return undefined; +class TableValuedExpandingLookup extends BasePreparedExpandingLookup { + constructor( + resultSetIndex: number, + readonly read: (request: RequestParameters) => SqliteParameterValue[][] + ) { + super(resultSetIndex); } - expandingLookupSync(stage: number, index: number): ParameterValueWithRow[][] | undefined { - const lookup = this.evaluators.lookupStages[stage][index]; - if (lookup.type == 'table_valued') { - // We can evaluate this table-valued function already. - const resultSetMarker = Symbol(); - const values = lookup.read(this.input.request).map((values, rowid) => { - const directOrigin: VirtualSourceRow = { resultSet: resultSetMarker, row: rowid }; - const provenance: [VirtualSourceRow] = [directOrigin]; - return values.map((value) => ({ value, provenance, directOrigin }) satisfies ParameterValueWithRow); - }); - - this.evaluators.lookupStages[stage][index] = { type: 'cached', values }; - return values; - } else if (lookup.type == 'cached') { - return lookup.values; - } - - return undefined; + override clone(): TableValuedExpandingLookup { + const lookup = new TableValuedExpandingLookup(this.resultSetIndex, this.read); + lookup.wasResolved = this.wasResolved; + return lookup; } } -class FullInstantiator extends PartialInstantiator { - resolveInstantiation(params: PreparedParameterValue[]): ParameterValueWithRow[][] { - const resolved = this.tryResolveInstantiation(params); - if (resolved == null) { - throw new Error('internal error: Should have been able to resolve instantiation after instantiating stages.'); - } - - return resolved; +class ParameterIndexExpandingLookup extends BasePreparedExpandingLookup { + constructor( + resultSetIndex: number, + readonly lookup: ParameterIndexLookupCreator, + readonly instantiation: PreparedParameterValue[] + ) { + super(resultSetIndex); } - async expandingLookup(stage: number, index: number): Promise { - const lookup = this.evaluators.lookupStages[stage][index]; - if (lookup.type == 'parameter') { - const scope = this.input.hydrationState.getParameterIndexLookupScope(lookup.lookup); - const resolvedLookup = lookup.lookup as PreparedParameterIndexLookupCreator; - - interface PendingLookup { - lookup: ScopedParameterLookup; - provenancePaths: VirtualSourceRow[][]; - resultSet: symbol; - } - - // It's possible that we'll have the same logical lookup with multiple provenance values. For instance, if the - // outputs of another lookup with two columns (where only one column is an input to this lookup) are passed into - // this, we can have two lookups with identical keys but different provenances. This hash map de-duplicates keys. - const pendingLookups = new HashMap( - FullInstantiator.parameterArrayEquality - ); - - for (const values of this.resolveInputs(lookup.instantiation)) { - const provenance: VirtualSourceRow[] = []; - for (const value of values) { - provenance.push(...value.provenance); - } - - const directValues = withoutProvenance(values); - pendingLookups.setOrUpdate(directValues, (old) => { - if (old == null) { - return { - lookup: ScopedParameterLookup.normalized(scope, UnscopedParameterLookup.normalized(directValues)), - provenancePaths: [provenance], - resultSet: Symbol(`lookup ${stage}.${index}`) - }; - } else { - old.provenancePaths.push(provenance); - return old; - } - }); - } - - const lookupsToProvenance = new Map(); - for (const [_, pending] of pendingLookups.entries) { - lookupsToProvenance.set(pending.lookup, pending); - } - - const outputs = await this.input.source.getParameterSets( - [...lookupsToProvenance.keys()], - `Stream ${this.evaluators.stream.name} evaluating parameter on ${resolvedLookup.sourceTable.tablePattern}` - ); - - const values = outputs.flatMap(({ lookup, rows }) => { - const { provenancePaths, resultSet } = lookupsToProvenance.get(lookup)!; - return provenancePaths.flatMap((origin, provenanceIndex) => { - return rows.map((row, rowid) => { - const length = Object.entries(row).length; - const asArray: ParameterValueWithRow[] = []; - - for (let i = 0; i < length; i++) { - // Stream parameters generate an output row like {0: , 1: , ...}. - const value = row[i.toString()] as SqliteParameterValue; - - // All paths share one result set because the lookup was deduplicated. Include the path index in the row - // identity because the same output rows are instantiated once for every path. - const directOrigin = length > 1 ? { resultSet, row: provenanceIndex * rows.length + rowid } : undefined; - - asArray.push({ - value, - // Note: Not tracking provenance for parameters with just a single output is purely a performance - // optimization. If there's just a single value, we don't need to correlate it with other columns in the - // row. Not adding provenance saves some work in mergeValueCombinations. - provenance: directOrigin != null ? [...origin, directOrigin] : origin, - directOrigin - }); - } - return asArray; - }); - }); - }); - - this.evaluators.lookupStages[stage][index] = { type: 'cached', values }; - return values; - } - - const other = this.expandingLookupSync(stage, index); - if (other == null) { - throw new Error('internal error: Unable to resolve non-parameter lookup synchronously?'); - } - return other; + override clone(): ParameterIndexExpandingLookup { + const lookup = new ParameterIndexExpandingLookup(this.resultSetIndex, this.lookup, this.instantiation); + lookup.wasResolved = this.wasResolved; + return lookup; } - - private static readonly parameterArrayEquality = listEquality(StableHasher.parameterValueEquality); } -export type PreparedExpandingLookup = - | { type: 'parameter'; lookup: ParameterIndexLookupCreator; instantiation: PreparedParameterValue[] } - | { type: 'table_valued'; read(request: RequestParameters): SqliteParameterValue[][] } - | { type: 'cached'; values: ParameterValueWithRow[][] }; - /** * A {@link plan.ParameterValue} that can be evaluated against request parameters. * - * Additionally, this includes the `cached` variant which allows partially instantiating parameters. + * Additionally, this includes the `static` variant which allows partially instantiating parameters. */ -export type PreparedParameterValue = - | { type: 'request'; read(request: RequestParameters): SqliteValue } - | { type: 'lookup'; lookup: { stage: number; index: number }; resultIndex: number } - | { type: 'intersection'; values: PreparedParameterValue[] } - | { type: 'cached'; values: ParameterValueWithRow[] }; +export type PreparedParameterValue = RequestParameterValue | LookupParameterValue; -interface ParameterValueWithRow { - value: SqliteParameterValue; +class RequestParameterValue { + resolved: SqliteParameterValue | undefined; - /** - * Information on how this value was resolved. - * - * We track how a parameter value was resolved to be able to merge parameters correctly. A Sync Stream with multiple - * independent parameters generates their cartesian product as buckets. When multiple parameters are resolved from the - * same row though, we can't use the full cartesian product. As an example, consider this stream: - * - * ```sql - * SELECT products.* FROM products, stores - * WHERE stores.name = products.store_name - * AND stores.region = products.region - * AND stores.id = subscription.parameter('store'); - * ``` - * - * Here, the bucket shape consists of two parameters (`store_name` and `region`). But since they're derived from the - * same `stores` row, we can't combine them freely. For each parameter, `store_name` and `region` must come from the - * same row. Here, the `provenance` would have a single entry and both parameters would have the same - * {@link VirtualSourceRow.resultSet}. - * - * For static values, such as constants or scalar values derived from request parameters, this array is empty. It's - * also possible for this to contain more than one entry, though: - * - * 1. For intersection values, we track the provenance of all input values. - * 2. It's possible to nest parameters. For instance, in the query `SELECT data.* FROM data, a, b WHERE a.a = data.a - * AND b.b = a.b AND data.c = b.c`, we have to parameters (`a` and `c`). `a` can be resolved from a lookup in - * table `a`, but we need to go through a second lookup to resolve `c`. Here, we can only combine value `c` with - * value `a` if the two were derived from the same row in `a`. So, the row `b` derived through `a` would include - * provenance elements of row `a` here. - */ - provenance: VirtualSourceRow[]; - // If set, must be contained in provenance - directOrigin?: VirtualSourceRow; -} + constructor(private readonly read: (request: RequestParameters) => SqliteValue) {} -interface VirtualSourceRow { - /** - * An opaque identifier for the result set this row was derived from. - */ - resultSet: symbol; - /** - * A number uniquely identifying this row in its result set. - */ - row: number; + resolveWith({ request }: PartialInstantiationInput): SqliteParameterValue { + if (this.resolved) return this.resolved; + + const value = this.read(request); + if (isValidParameterValue(value)) { + return (this.resolved = value); + } else { + throw uninstantiableException; + } + } + + clone(): RequestParameterValue { + const clone = new RequestParameterValue(this.read); + clone.resolved = this.resolved; + return clone; + } } -function withoutProvenance(source: ParameterValueWithRow[]): SqliteParameterValue[] { - return source.map(({ value }) => value); +class LookupParameterValue implements ResultSetColumn { + constructor( + readonly lookup: PreparedExpandingLookup, + readonly outputIndex: number + ) {} + + clone(): LookupParameterValue { + return new LookupParameterValue(this.lookup, this.outputIndex); + } } export interface PartialInstantiationInput { @@ -691,64 +456,3 @@ function* filterParameterRows(rows: SqliteValue[][]): Generator { - // Partial backtracking results, the current instantiation is fixed for 0..nextParameter in generateCombinations. - const partialResults = new Array(valuesByParameter.length); - // A map from result sets to rows used in the partial instantiation. - const usedRows = new Map(); - - function installRowIfNoConflict(value: ParameterValueWithRow): [boolean, symbol[]] { - const addedResultSets: symbol[] = []; - - for (const origin of value.provenance) { - const { resultSet, row } = origin; - const existingRow = usedRows.get(resultSet); - if (existingRow === undefined) { - addedResultSets.push(resultSet); - usedRows.set(resultSet, row); - } else if (existingRow == row) { - continue; - } else { - // The current instantiation already contains a value from the same result set but derived from a different - // row. So we must ignore this parameter value. - return [false, addedResultSets]; - } - } - - return [true, addedResultSets]; - } - - function uninstallResultSets(resultSets: symbol[]) { - for (const rs of resultSets) { - usedRows.delete(rs); - } - } - - function* generateCombinations(nextParameter: number): Generator { - if (nextParameter >= valuesByParameter.length) { - yield [...partialResults]; - return; - } - - const availableValues = valuesByParameter[nextParameter]; - for (const available of availableValues) { - const [canUse, addedResultSets] = installRowIfNoConflict(available); - if (canUse) { - partialResults[nextParameter] = available; - yield* generateCombinations(nextParameter + 1); - } - - uninstallResultSets(addedResultSets); - } - } - - yield* generateCombinations(0); -} diff --git a/packages/sync-rules/src/sync_plan/evaluator/result_set.ts b/packages/sync-rules/src/sync_plan/evaluator/result_set.ts new file mode 100644 index 000000000..781ecd080 --- /dev/null +++ b/packages/sync-rules/src/sync_plan/evaluator/result_set.ts @@ -0,0 +1,177 @@ +import { HashMap, listEquality, StableHasher } from '../../compiler/equality.js'; +import { SqliteParameterValue } from '../../types.js'; + +/** + * A mutable result set of parameter results. + * + * This is used to represent parameter results when resolving buckets: Each expanding lookup is joined onto a pending + * result set until all lookups have been applied. Once all result sets have been added, bucket parameters can be read + * by reading columns in each row. + */ +export class ResultSet { + #totalLookups: number; + #rows: ResultSetRow[]; + + constructor(totalLookups: number) { + this.#totalLookups = totalLookups; + const initialRow = new Array(totalLookups); + initialRow.fill(undefined); + this.#rows = [initialRow]; + } + + get length(): number { + return this.#rows.length; + } + + clone(): ResultSet { + const rs = new ResultSet(this.#totalLookups); + rs.#rows.splice(0, 1); // Remove the initial unit row + + for (const row of this.#rows) { + rs.#rows.push([...row]); + } + return rs; + } + + /** + * Extracts unique values by looking values for each column in this result set. + * + * For each unique projection row, also returns all source rows the value was derived from. + */ + *projectUnique(columns: ResultSetColumn[]): Iterable { + for (const { group, first } of this.#groupBy(columns, (values) => values)) { + if (first) { + yield group; + } + } + } + + /** + * Adds a new result set by forming the cartesian product with the given values. + */ + multiply(resultSetIndex: number, rows: SqliteParameterValue[][]) { + if (rows.length === 0) { + this.#rows = []; + } + + const originalLength = this.#rows.length; + for (let i = 0; i < originalLength; i++) { + this.#multiplyAtRow(resultSetIndex, i, rows); + } + } + + async joinAsync( + keys: ResultSetColumn[], + resultSetIndex: number, + performLookup: (lookups: AsyncJoinLookup[]) => Promise + ) { + const lookupsByRow: AsyncJoinLookup[] = []; + const uniqueLookups: AsyncJoinLookup[] = []; + + for (const { group, first } of this.#groupBy(keys, (values) => ({ inputs: values, foundRows: [] }))) { + if (first) uniqueLookups.push(group); + lookupsByRow.push(group); + } + + await performLookup(uniqueLookups); + + const deletedRows: number[] = []; + const originalLength = this.#rows.length; + for (let i = 0; i < originalLength; i++) { + const lookup = lookupsByRow[i]; + if (lookup.foundRows.length > 0) { + this.#multiplyAtRow(resultSetIndex, i, lookup.foundRows); + } else { + // The row has no matching join partner, so remove it. We can't split it immediately because #multiplyAtRow is + // still iterating through rows. + deletedRows.push(i); + } + } + + let offset = 0; + for (const toDelete of deletedRows) { + this.#rows.splice(toDelete - offset, 1); + offset++; + } + } + + #multiplyAtRow(resultSetIndex: number, rowIndex: number, rows: SqliteParameterValue[][]) { + // Add first element of product to existing row, remaining as new rows. + const row = this.#rows[rowIndex]; + row[resultSetIndex] = rows[0]; + + for (let j = 1; j < rows.length; j++) { + const copy = Array.from(row); + copy[resultSetIndex] = rows[j]; + this.#rows.push(copy); + } + } + + *#groupBy(columns: ResultSetColumn[], generateGroup: (values: SqliteParameterValue[]) => T) { + const originalLength = this.#rows.length; + + if (columns.length === 1) { + // Fast path, we can use native sets. + const [column] = columns; + const foundValues = new Map(); + + for (let i = 0; i < originalLength; i++) { + const row = this.#rows[i]; + const value = lookupInRow(row, column); + const existingGroup = foundValues.get(value); + + if (existingGroup != null) { + yield { index: i, group: existingGroup, first: false }; + } else { + const group = generateGroup([value]); + foundValues.set(value, group); + yield { index: i, group, first: true }; + } + } + } else { + const foundValues = new HashMap(parameterArrayEquality); + + for (let i = 0; i < originalLength; i++) { + const row = this.#rows[i]; + const values = columns.map((c) => lookupInRow(row, c)); + + let isFirst = false; + const group = foundValues.putIfAbsent(values, () => { + isFirst = true; + return generateGroup(values); + }); + + yield { index: i, group, first: isFirst }; + } + } + } +} + +export interface ResultSetElement { + resultSetIndex: number; +} + +export interface ResultSetColumn { + lookup: ResultSetElement; + outputIndex: number; +} + +export interface AsyncJoinLookup { + inputs: SqliteParameterValue[]; + foundRows: SqliteParameterValue[][]; +} + +/** + * A row in a result set. + * + * While this is semantically a list of columns, that representation would require a lot of copying on each join. + * So, we represent each lookup result as an array of values (that we can re-use when we create new rows for joins). + * Result sets that have not yet been processed are represented as undefined. + */ +type ResultSetRow = (SqliteParameterValue[] | undefined)[]; + +function lookupInRow(row: ResultSetRow, column: ResultSetColumn): SqliteParameterValue { + return row[column.lookup.resultSetIndex]![column.outputIndex]; +} + +const parameterArrayEquality = listEquality(StableHasher.parameterValueEquality); diff --git a/packages/sync-rules/test/src/sync_plan/evaluator/result_set.test.ts b/packages/sync-rules/test/src/sync_plan/evaluator/result_set.test.ts new file mode 100644 index 000000000..aca467927 --- /dev/null +++ b/packages/sync-rules/test/src/sync_plan/evaluator/result_set.test.ts @@ -0,0 +1,179 @@ +import { describe, expect, test } from 'vitest'; +import { ResultSet, ResultSetElement } from '../../../../src/sync_plan/evaluator/result_set.js'; +import { SqliteParameterValue } from '../../../../src/types.js'; + +describe('ResultSet', () => { + test('unit result set', () => { + const empty = new ResultSet(0); + + expect(empty.length).toStrictEqual(1); + expect([...empty.projectUnique([])]).toStrictEqual([[]]); + }); + + describe('projectUnique', () => { + const rs = new ResultSet(1); + rs.multiply(0, [ + ['a1', 'b1'], + ['a2', 'b2'], + ['a2', 'b1'], + ['a2', 'b2'] + ]); + + test('single column', () => { + expect([...rs.projectUnique([{ lookup: element(0), outputIndex: 0 }])]).toStrictEqual([['a1'], ['a2']]); + + expect([...rs.projectUnique([{ lookup: element(0), outputIndex: 1 }])]).toStrictEqual([['b1'], ['b2']]); + }); + + test('multiple columns', () => { + expect([ + ...rs.projectUnique([ + { lookup: element(0), outputIndex: 0 }, + { lookup: element(0), outputIndex: 1 } + ]) + ]).toStrictEqual([ + ['a1', 'b1'], + ['a2', 'b2'], + ['a2', 'b1'] + ]); + }); + }); + + describe('multiply', () => { + test('empty', () => { + const rs = new ResultSet(1); + expect(rs.length).toStrictEqual(1); + + rs.multiply(0, []); + expect(rs.length).toStrictEqual(0); + }); + + test('is cartesian product', () => { + const rs = new ResultSet(2); + + rs.multiply(0, [['a'], ['b']]); + rs.multiply(1, [[0], [1]]); + expect([ + ...rs.projectUnique([ + { + lookup: element(0), + outputIndex: 0 + }, + { + lookup: element(1), + outputIndex: 0 + } + ]) + ]).toStrictEqual([ + ['a', 0], + ['b', 0], + ['a', 1], + ['b', 1] + ]); + }); + }); + + describe('joinAsync', () => { + const col0 = { lookup: element(0), outputIndex: 0 }; + const col1 = { lookup: element(1), outputIndex: 0 }; + + test('expands each row with matching values, deduplicating lookups', async () => { + const rs = new ResultSet(2); + rs.multiply(0, [['a'], ['b'], ['a']]); + + const matches: Record = { + a: [[1], [2]], + b: [[3]] + }; + + let lookupCount = 0; + await rs.joinAsync([col0], 1, async (lookups) => { + // The lookup for 'a' is only performed once, even though it's shared by two rows. + expect(lookups.length).toStrictEqual(2); + lookupCount++; + + for (const lookup of lookups) { + lookup.foundRows.push(...matches[lookup.inputs[0] as string]); + } + }); + + expect(lookupCount).toStrictEqual(1); + expect([...rs.projectUnique([col0, col1])]).toStrictEqual([ + ['a', 1], + ['b', 3], + ['a', 2] + ]); + }); + + test('removes rows without a matching join partner', async () => { + const rs = new ResultSet(2); + rs.multiply(0, [['a'], ['b'], ['c'], ['a']]); + + const matches: Record = { + a: [[10]], + b: [], + c: [[30], [31]] + }; + + await rs.joinAsync([col0], 1, async (lookups) => { + for (const lookup of lookups) { + lookup.foundRows.push(...matches[lookup.inputs[0] as string]); + } + }); + + expect(rs.length).toStrictEqual(4); + expect([...rs.projectUnique([col0, col1])]).toStrictEqual([ + ['a', 10], + ['c', 30], + ['c', 31] + ]); + }); + + test('removing all rows results in an empty result set', async () => { + const rs = new ResultSet(2); + rs.multiply(0, [['a'], ['b']]); + + await rs.joinAsync([col0], 1, async () => { + // Leave foundRows empty for every lookup. + }); + + expect(rs.length).toStrictEqual(0); + expect([...rs.projectUnique([col0, col1])]).toStrictEqual([]); + }); + + test('supports composite join keys', async () => { + const rs = new ResultSet(2); + rs.multiply(0, [ + [1, 'x'], + [1, 'y'], + [2, 'x'] + ]); + + const colKey0 = { lookup: element(0), outputIndex: 0 }; + const colKey1 = { lookup: element(0), outputIndex: 1 }; + + const matches: Record = { + '1,x': [[100]], + '1,y': [[200]], + '2,x': [[300], [301]] + }; + + await rs.joinAsync([colKey0, colKey1], 1, async (lookups) => { + for (const lookup of lookups) { + lookup.foundRows.push(...matches[lookup.inputs.join(',')]); + } + }); + + expect([...rs.projectUnique([colKey0, colKey1, col1])]).toStrictEqual([ + [1, 'x', 100], + [1, 'y', 200], + [2, 'x', 300], + [2, 'x', 301] + ]); + }); + }); +}); + +function element(index: number): ResultSetElement { + return { resultSetIndex: index }; +} From 46d937444d54422af33982cdba5c4b690f2ef8ab Mon Sep 17 00:00:00 2001 From: Simon Binder Date: Tue, 1 Sep 2026 15:25:02 +0200 Subject: [PATCH 02/15] Fix most evaluator tests --- .../src/sync_plan/evaluator/parameter_evaluator.ts | 3 +-- .../test/src/sync_plan/evaluator/evaluator.test.ts | 8 ++------ 2 files changed, 3 insertions(+), 8 deletions(-) diff --git a/packages/sync-rules/src/sync_plan/evaluator/parameter_evaluator.ts b/packages/sync-rules/src/sync_plan/evaluator/parameter_evaluator.ts index dffe7ca68..a89f6c5fa 100644 --- a/packages/sync-rules/src/sync_plan/evaluator/parameter_evaluator.ts +++ b/packages/sync-rules/src/sync_plan/evaluator/parameter_evaluator.ts @@ -276,12 +276,10 @@ export class RequestParameterEvaluators { } for (const stage of lookupStages) { - const stageIndex = mappedStages.length; const mappedStage: PreparedExpandingLookup[] = []; mappedStages.push(mappedStage); for (const lookup of stage) { - const index = mappedStage.length; let resolved: PreparedExpandingLookup; if (lookup.type == 'parameter') { @@ -316,6 +314,7 @@ export class RequestParameterEvaluators { } lookupToStage.set(lookup, resolved); + mappedStage.push(resolved); } } diff --git a/packages/sync-rules/test/src/sync_plan/evaluator/evaluator.test.ts b/packages/sync-rules/test/src/sync_plan/evaluator/evaluator.test.ts index e1ee6837c..27e05b673 100644 --- a/packages/sync-rules/test/src/sync_plan/evaluator/evaluator.test.ts +++ b/packages/sync-rules/test/src/sync_plan/evaluator/evaluator.test.ts @@ -846,11 +846,7 @@ streams: // Duplicates do not need to be removed here, but they must not make lookup // columns independent and create impossible pairs like [2, "A"]. - expect(dynamicBuckets.map((bucket) => bucket.bucket)).toStrictEqual([ - 'stream|0[1,"A"]', - 'stream|0[1,"A"]', - 'stream|0[2,"B"]' - ]); + expect(dynamicBuckets.map((bucket) => bucket.bucket)).toStrictEqual(['stream|0[1,"A"]', 'stream|0[2,"B"]']); }); syncTest('preserves correlation across bigint lookup output columns', async ({ sync }) => { @@ -1103,8 +1099,8 @@ streams: expect(querier.staticBuckets.map((e) => e.bucket)).toStrictEqual([ 'stream|0["a1","b1"]', - 'stream|0["a1","b2"]', 'stream|0["a2","b1"]', + 'stream|0["a1","b2"]', 'stream|0["a2","b2"]' ]); }); From 0454ab0857cd7e88faab656c2b8b9d46892f3b9d Mon Sep 17 00:00:00 2001 From: Simon Binder Date: Tue, 1 Sep 2026 16:50:12 +0200 Subject: [PATCH 03/15] Support intersections --- .../evaluator/parameter_evaluator.ts | 175 +++++++++++++++--- .../src/sync_plan/evaluator/result_set.ts | 24 +++ .../src/sync_plan/evaluator/evaluator.test.ts | 26 +++ .../sync_plan/evaluator/result_set.test.ts | 28 +++ 4 files changed, 228 insertions(+), 25 deletions(-) diff --git a/packages/sync-rules/src/sync_plan/evaluator/parameter_evaluator.ts b/packages/sync-rules/src/sync_plan/evaluator/parameter_evaluator.ts index a89f6c5fa..23ebab82d 100644 --- a/packages/sync-rules/src/sync_plan/evaluator/parameter_evaluator.ts +++ b/packages/sync-rules/src/sync_plan/evaluator/parameter_evaluator.ts @@ -59,11 +59,11 @@ export class RequestParameterEvaluators { /** * Pending lookup stages, or their cached outputs. */ - readonly lookupStages: PreparedExpandingLookup[][], + private readonly lookupStages: LookupStage[], /** * Pending parameter values, or their cached outputs. */ - readonly parameterValues: PreparedParameterValue[], + private readonly parameterValues: PreparedParameterValue[], /** * The materialized result set into which @@ -82,8 +82,19 @@ export class RequestParameterEvaluators { * instead of re-evaluating them on every parameter lookup change. */ clone(): RequestParameterEvaluators { - const copiedStages = this.lookupStages.map((s) => s.map((e) => e.clone())); - const outputValues = this.parameterValues.map((v) => v.clone()); + const clonedParameters = new Map(); + + function cloneParameter(original: PreparedParameterValue) { + const existing = clonedParameters.get(original); + if (existing != null) return existing; + + const clone = original.clone(); + clonedParameters.set(original, clone); + return clone; + } + + const copiedStages = this.lookupStages.map((s) => s.clone(cloneParameter)); + const outputValues = this.parameterValues.map(cloneParameter); return new RequestParameterEvaluators(this.stream, copiedStages, outputValues, this.resultSet.clone()); } @@ -100,7 +111,9 @@ export class RequestParameterEvaluators { try { // At this point, we can resolve table-valued lookups and parameter values based only on request data. for (const stage of this.lookupStages) { - for (const element of stage) { + let needsParameterLookups = false; + + for (const element of stage.lookups) { if (element instanceof TableValuedExpandingLookup) { const outputs = element.read(input.request); element.wasResolved = true; @@ -108,9 +121,21 @@ export class RequestParameterEvaluators { this.#checkInstantiable(); } else { - for (const instantiation of element.instantiation) { - if (instantiation instanceof RequestParameterValue) instantiation.resolveWith(input); - } + needsParameterLookups = true; + } + } + + for (const instantiation of stage.inputParameters()) { + if (instantiation instanceof RequestParameterValue) { + instantiation.resolveWith(input); + } else { + needsParameterLookups = true; + } + } + + if (!needsParameterLookups) { + for (const intersection of stage.intersections) { + this.#applyIntersectionConstraint(intersection); } } } @@ -134,12 +159,16 @@ export class RequestParameterEvaluators { */ async instantiate(input: InstantiationInput): Promise { try { - for (const stage of this.lookupStages) { - for (const lookup of stage) { + for (const { lookups, intersections } of this.lookupStages) { + for (const lookup of lookups) { if (lookup instanceof ParameterIndexExpandingLookup) { await this.#instantiateLookup(lookup, input); } } + + for (const intersection of intersections) { + this.#applyIntersectionConstraint(intersection); + } } return this.#readParameters()!; @@ -155,11 +184,13 @@ export class RequestParameterEvaluators { } #readParameters(): SqliteParameterValue[][] | undefined { - for (const stage of this.lookupStages) { - for (const element of stage) { - if (!element.wasResolved) { - return undefined; - } + for (const { intersections, lookups } of this.lookupStages) { + for (const intersection of intersections) { + if (!intersection.wasApplied) return undefined; + } + + for (const element of lookups) { + if (!element.wasResolved) return undefined; } } @@ -183,12 +214,34 @@ export class RequestParameterEvaluators { if (v instanceof LookupParameterValue) { return row[lookupIndex++]; } else { - if (!v.resolved) throw new Error('Expected request values to be resolved here'); - return v.resolved; + return v.requireResolved(); } }); } + #applyIntersectionConstraint(constraint: RequiredIntersection) { + // If any parameter of the intersection is a scalar value derived from a request, that value. + let knownValue: SqliteParameterValue | undefined; + const intersection: ResultSetColumn[] = []; + + for (const value of constraint.values) { + if (value instanceof RequestParameterValue) { + const evaluated = value.requireResolved(); + if (knownValue !== undefined && evaluated != knownValue) { + throw uninstantiableException; + } + + knownValue = evaluated; + } else { + intersection.push(value); + } + } + + this.resultSet.formIntersection(intersection, knownValue); + constraint.wasApplied = true; + this.#checkInstantiable(); + } + async #instantiateLookup(lookup: ParameterIndexExpandingLookup, input: InstantiationInput) { const scope = input.hydrationState.getParameterIndexLookupScope(lookup.lookup); const resolvedLookup = lookup.lookup as PreparedParameterIndexLookupCreator; @@ -249,7 +302,7 @@ export class RequestParameterEvaluators { input: StreamInput, engine: ScalarExpressionEngine ) { - const mappedStages: PreparedExpandingLookup[][] = []; + const mappedStages: LookupStage[] = []; let amountOfLookups = 0; const lookupToStage = new Map(); @@ -267,7 +320,19 @@ export class RequestParameterEvaluators { const lookup = lookupToStage.get(value.lookup)!; return new LookupParameterValue(lookup, value.resultIndex); } else { - throw new Error('TODO: intersection'); + const intersectionInputs = mapParameterValues(value.values); + + if (mappedStages.length > 0) { + mappedStages[mappedStages.length - 1].intersections.push({ values: intersectionInputs, wasApplied: false }); + } else { + // Intersection in first stage, e.g. for request parameters. Add a stage just for this. + const stage = new LookupStage([], [{ values: intersectionInputs, wasApplied: false }]); + mappedStages.push(stage); + } + + // Non-intersecting rows will be pruned from the result set or, for scalar parameter values, mark the querier + // as uninstantiable. So, we can replace the intersection value with any inner value. + return intersectionInputs[0]; } } @@ -276,8 +341,7 @@ export class RequestParameterEvaluators { } for (const stage of lookupStages) { - const mappedStage: PreparedExpandingLookup[] = []; - mappedStages.push(mappedStage); + const mappedStage = new LookupStage([], []); for (const lookup of stage) { let resolved: PreparedExpandingLookup; @@ -314,8 +378,10 @@ export class RequestParameterEvaluators { } lookupToStage.set(lookup, resolved); - mappedStage.push(resolved); + mappedStage.lookups.push(resolved); } + + mappedStages.push(mappedStage); } const rs = new ResultSet(amountOfLookups); @@ -332,12 +398,14 @@ const uninstantiableException = Symbol.for('uninstantiable'); export type PreparedExpandingLookup = TableValuedExpandingLookup | ParameterIndexExpandingLookup; +type CloneParameter = (original: PreparedParameterValue) => PreparedParameterValue; + abstract class BasePreparedExpandingLookup implements ResultSetElement { wasResolved = false; constructor(readonly resultSetIndex: number) {} - abstract clone(): BasePreparedExpandingLookup; + abstract clone(parameters: CloneParameter): BasePreparedExpandingLookup; } class TableValuedExpandingLookup extends BasePreparedExpandingLookup { @@ -364,13 +432,56 @@ class ParameterIndexExpandingLookup extends BasePreparedExpandingLookup { super(resultSetIndex); } - override clone(): ParameterIndexExpandingLookup { - const lookup = new ParameterIndexExpandingLookup(this.resultSetIndex, this.lookup, this.instantiation); + override clone(cloneParameter: CloneParameter): ParameterIndexExpandingLookup { + const lookup = new ParameterIndexExpandingLookup( + this.resultSetIndex, + this.lookup, + this.instantiation.map(cloneParameter) + ); lookup.wasResolved = this.wasResolved; return lookup; } } +class LookupStage { + constructor( + /** + * Lookups that only have dependencies on prior stages. + */ + readonly lookups: PreparedExpandingLookup[], + /** + * A list of constraints enforcing that specific columns must have equal values. + * + * These constraints are evaluated after lookups, and may reference lookups in this stage. + */ + readonly intersections: RequiredIntersection[] + ) {} + + *inputParameters() { + for (const intersection of this.intersections) { + yield* intersection.values; + } + + for (const lookup of this.lookups) { + if (lookup instanceof ParameterIndexExpandingLookup) { + yield* lookup.instantiation; + } + } + } + + clone(cloneParameter: CloneParameter): LookupStage { + return new LookupStage( + this.lookups.map((l) => l.clone(cloneParameter)), + this.intersections.map(({ values, wasApplied }) => ({ values: values.map(cloneParameter), wasApplied })) + ); + } +} + +interface RequiredIntersection { + values: PreparedParameterValue[]; + wasApplied: boolean; +} + /** * A {@link plan.ParameterValue} that can be evaluated against request parameters. * @@ -383,6 +494,11 @@ class RequestParameterValue { constructor(private readonly read: (request: RequestParameters) => SqliteValue) {} + requireResolved() { + if (!this.resolved) throw new Error('Expected request values to be resolved here'); + return this.resolved; + } + resolveWith({ request }: PartialInstantiationInput): SqliteParameterValue { if (this.resolved) return this.resolved; @@ -455,3 +571,12 @@ function* filterParameterRows(rows: SqliteValue[][]): Generator e.bucket)).toStrictEqual(['stream|0["user"]']); }); + syncTest('intersection of request data', ({ sync }) => { + const desc = sync.prepareSyncStreams(` +config: + edition: 3 + +streams: + stream: + auto_subscribe: true + query: SELECT * FROM issues WHERE a = auth.parameter('x') AND a = auth.parameter('y') +`); + + function queryWith(x: string, y: string) { + const { querier, errors } = desc.getBucketParameterQuerier({ + globalParameters: requestParameters({ sub: 'user', x, y }), + hasDefaultStreams: true, + streams: {} + }); + expect(errors).toStrictEqual([]); + expect(querier.hasDynamicBuckets).toStrictEqual(false); + return querier.staticBuckets.map((e) => e.bucket); + } + + expect(queryWith('p1', 'p2')).toStrictEqual([]); + expect(queryWith('p1', 'p1')).toStrictEqual(['stream|0["p1"]']); + }); + syncTest('parameter lookups', async ({ sync }) => { const desc = sync.prepareSyncStreams(` config: diff --git a/packages/sync-rules/test/src/sync_plan/evaluator/result_set.test.ts b/packages/sync-rules/test/src/sync_plan/evaluator/result_set.test.ts index aca467927..5a9bf8b78 100644 --- a/packages/sync-rules/test/src/sync_plan/evaluator/result_set.test.ts +++ b/packages/sync-rules/test/src/sync_plan/evaluator/result_set.test.ts @@ -73,6 +73,34 @@ describe('ResultSet', () => { }); }); + describe('formIntersection', () => { + const col0 = { lookup: element(0), outputIndex: 0 }; + const col1 = { lookup: element(1), outputIndex: 0 }; + + test('with a fixed value, removes rows where any column differs from it', () => { + const rs = new ResultSet(2); + rs.multiply(0, [['a'], ['b']]); + rs.multiply(1, [['a'], ['b']]); + + rs.formIntersection([col0, col1], 'a'); + + expect([...rs.projectUnique([col0, col1])]).toStrictEqual([['a', 'a']]); + }); + + test('without a fixed value, removes rows where the columns differ from each other', () => { + const rs = new ResultSet(2); + rs.multiply(0, [['a'], ['b']]); + rs.multiply(1, [['a'], ['b']]); + + rs.formIntersection([col0, col1]); + + expect([...rs.projectUnique([col0, col1])]).toStrictEqual([ + ['a', 'a'], + ['b', 'b'] + ]); + }); + }); + describe('joinAsync', () => { const col0 = { lookup: element(0), outputIndex: 0 }; const col1 = { lookup: element(1), outputIndex: 0 }; From a17c768cbfb7b41920dd8b5f58f79c156f627ba8 Mon Sep 17 00:00:00 2001 From: Simon Binder Date: Tue, 1 Sep 2026 17:38:32 +0200 Subject: [PATCH 04/15] Test empty joins --- .../evaluator/parameter_evaluator.ts | 7 ++- .../src/sync_plan/evaluator/result_set.ts | 61 +++++++++++-------- .../sync_plan/evaluator/result_set.test.ts | 16 +++++ 3 files changed, 56 insertions(+), 28 deletions(-) diff --git a/packages/sync-rules/src/sync_plan/evaluator/parameter_evaluator.ts b/packages/sync-rules/src/sync_plan/evaluator/parameter_evaluator.ts index 23ebab82d..bd7c08746 100644 --- a/packages/sync-rules/src/sync_plan/evaluator/parameter_evaluator.ts +++ b/packages/sync-rules/src/sync_plan/evaluator/parameter_evaluator.ts @@ -55,7 +55,7 @@ import { AsyncJoinLookup, ResultSet, ResultSetColumn, ResultSetElement } from '. */ export class RequestParameterEvaluators { private constructor( - readonly stream: plan.StreamOptions, + private readonly stream: plan.StreamOptions, /** * Pending lookup stages, or their cached outputs. */ @@ -66,9 +66,10 @@ export class RequestParameterEvaluators { private readonly parameterValues: PreparedParameterValue[], /** - * The materialized result set into which + * The materialized result set containing lookup values. {@link parameterValues} are read from this result set as a + * final step. */ - readonly resultSet: ResultSet + private readonly resultSet: ResultSet ) {} /** diff --git a/packages/sync-rules/src/sync_plan/evaluator/result_set.ts b/packages/sync-rules/src/sync_plan/evaluator/result_set.ts index 1fc299672..c7747a57e 100644 --- a/packages/sync-rules/src/sync_plan/evaluator/result_set.ts +++ b/packages/sync-rules/src/sync_plan/evaluator/result_set.ts @@ -12,6 +12,9 @@ export class ResultSet { #totalLookups: number; #rows: ResultSetRow[]; + /** + * @param totalLookups - The total amount of lookups that will be joined to this result set. + */ constructor(totalLookups: number) { this.#totalLookups = totalLookups; const initialRow = new Array(totalLookups); @@ -35,8 +38,6 @@ export class ResultSet { /** * Extracts unique values by looking values for each column in this result set. - * - * For each unique projection row, also returns all source rows the value was derived from. */ *projectUnique(columns: ResultSetColumn[]): Iterable { for (const { group, first } of this.#groupBy(columns, (values) => values)) { @@ -61,29 +62,10 @@ export class ResultSet { } /** - * Removes rows where the given columns have different values. - * - * If a fixed value is passed, this also removes rows where any of the given columns has a different value. + * @param keys - Join keys that are already present in the result set. + * @param resultSetIndex - The index of the resl set being joined. + * @param performLookup - Adds resolved rows to each unique instantiation of join keys. */ - formIntersection(columns: ResultSetColumn[], fixedValue?: SqliteParameterValue) { - row: for (let i = 0; i < this.#rows.length; i++) { - const row = this.#rows[i]; - let requiredValue = fixedValue; - - for (const column of columns) { - const evaluated = lookupInRow(row, column); - if (requiredValue !== undefined && evaluated != requiredValue) { - // This row needs to be removed! - this.#rows.splice(i, 1); - i--; - continue row; - } - - requiredValue = evaluated; - } - } - } - async joinAsync( keys: ResultSetColumn[], resultSetIndex: number, @@ -119,6 +101,30 @@ export class ResultSet { } } + /** + * Removes rows where the given columns have different values. + * + * If a fixed value is passed, this also removes rows where any of the given columns has a different value. + */ + formIntersection(columns: ResultSetColumn[], fixedValue?: SqliteParameterValue) { + row: for (let i = 0; i < this.#rows.length; i++) { + const row = this.#rows[i]; + let requiredValue = fixedValue; + + for (const column of columns) { + const evaluated = lookupInRow(row, column); + if (requiredValue !== undefined && evaluated != requiredValue) { + // This row needs to be removed! + this.#rows.splice(i, 1); + i--; + continue row; + } + + requiredValue = evaluated; + } + } + } + #multiplyAtRow(resultSetIndex: number, rowIndex: number, rows: SqliteParameterValue[][]) { // Add first element of product to existing row, remaining as new rows. const row = this.#rows[rowIndex]; @@ -195,7 +201,12 @@ export interface AsyncJoinLookup { type ResultSetRow = (SqliteParameterValue[] | undefined)[]; function lookupInRow(row: ResultSetRow, column: ResultSetColumn): SqliteParameterValue { - return row[column.lookup.resultSetIndex]![column.outputIndex]; + const valuesForResultSet = row[column.lookup.resultSetIndex]; + if (valuesForResultSet === undefined) { + throw new Error('Tried to lookup values set before it was joined to result set'); + } + + return valuesForResultSet[column.outputIndex]; } const parameterArrayEquality = listEquality(StableHasher.parameterValueEquality); diff --git a/packages/sync-rules/test/src/sync_plan/evaluator/result_set.test.ts b/packages/sync-rules/test/src/sync_plan/evaluator/result_set.test.ts index 5a9bf8b78..27490cf55 100644 --- a/packages/sync-rules/test/src/sync_plan/evaluator/result_set.test.ts +++ b/packages/sync-rules/test/src/sync_plan/evaluator/result_set.test.ts @@ -199,6 +199,22 @@ describe('ResultSet', () => { [2, 'x', 301] ]); }); + + test('empty join key', async () => { + const rs = new ResultSet(2); + rs.multiply(0, [['a'], ['b']]); + + await rs.joinAsync([], 1, async (lookups) => { + expect(lookups).toMatchObject([{ inputs: [] }]); + + lookups[0].foundRows.push(['x']); + }); + + expect([...rs.projectUnique([col0, col1])]).toStrictEqual([ + ['a', 'x'], + ['b', 'x'] + ]); + }); }); }); From a31864aa229b2a5353919b4cd9f825b7a10b5922 Mon Sep 17 00:00:00 2001 From: Simon Binder Date: Tue, 1 Sep 2026 17:49:04 +0200 Subject: [PATCH 05/15] Fix for falsy resolved values --- .../src/sync_plan/evaluator/parameter_evaluator.ts | 12 ++++++------ 1 file changed, 6 insertions(+), 6 deletions(-) diff --git a/packages/sync-rules/src/sync_plan/evaluator/parameter_evaluator.ts b/packages/sync-rules/src/sync_plan/evaluator/parameter_evaluator.ts index bd7c08746..a1a9175f3 100644 --- a/packages/sync-rules/src/sync_plan/evaluator/parameter_evaluator.ts +++ b/packages/sync-rules/src/sync_plan/evaluator/parameter_evaluator.ts @@ -491,21 +491,21 @@ interface RequiredIntersection { export type PreparedParameterValue = RequestParameterValue | LookupParameterValue; class RequestParameterValue { - resolved: SqliteParameterValue | undefined; + #resolved: SqliteParameterValue | undefined; constructor(private readonly read: (request: RequestParameters) => SqliteValue) {} requireResolved() { - if (!this.resolved) throw new Error('Expected request values to be resolved here'); - return this.resolved; + if (this.#resolved === undefined) throw new Error('Expected request values to be resolved here'); + return this.#resolved; } resolveWith({ request }: PartialInstantiationInput): SqliteParameterValue { - if (this.resolved) return this.resolved; + if (this.#resolved !== undefined) return this.#resolved; const value = this.read(request); if (isValidParameterValue(value)) { - return (this.resolved = value); + return (this.#resolved = value); } else { throw uninstantiableException; } @@ -513,7 +513,7 @@ class RequestParameterValue { clone(): RequestParameterValue { const clone = new RequestParameterValue(this.read); - clone.resolved = this.resolved; + clone.#resolved = this.#resolved; return clone; } } From 69caf10ec27a1db3c0a0ff7511f7a3f23c65bc79 Mon Sep 17 00:00:00 2001 From: Simon Binder Date: Wed, 2 Sep 2026 09:40:38 +0200 Subject: [PATCH 06/15] Freeze rows before adding them --- packages/sync-rules/src/sync_plan/evaluator/result_set.ts | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/packages/sync-rules/src/sync_plan/evaluator/result_set.ts b/packages/sync-rules/src/sync_plan/evaluator/result_set.ts index c7747a57e..6c41cddf0 100644 --- a/packages/sync-rules/src/sync_plan/evaluator/result_set.ts +++ b/packages/sync-rules/src/sync_plan/evaluator/result_set.ts @@ -31,6 +31,7 @@ export class ResultSet { rs.#rows.splice(0, 1); // Remove the initial unit row for (const row of this.#rows) { + // We can shallow-clone rows, inner items are frozen once added into the result set. rs.#rows.push([...row]); } return rs; @@ -128,11 +129,11 @@ export class ResultSet { #multiplyAtRow(resultSetIndex: number, rowIndex: number, rows: SqliteParameterValue[][]) { // Add first element of product to existing row, remaining as new rows. const row = this.#rows[rowIndex]; - row[resultSetIndex] = rows[0]; + row[resultSetIndex] = Object.freeze(rows[0]); for (let j = 1; j < rows.length; j++) { const copy = Array.from(row); - copy[resultSetIndex] = rows[j]; + copy[resultSetIndex] = Object.freeze(rows[j]); this.#rows.push(copy); } } @@ -198,7 +199,7 @@ export interface AsyncJoinLookup { * So, we represent each lookup result as an array of values (that we can re-use when we create new rows for joins). * Result sets that have not yet been processed are represented as undefined. */ -type ResultSetRow = (SqliteParameterValue[] | undefined)[]; +type ResultSetRow = (ReadonlyArray | undefined)[]; function lookupInRow(row: ResultSetRow, column: ResultSetColumn): SqliteParameterValue { const valuesForResultSet = row[column.lookup.resultSetIndex]; From b6ff8f3d8c4e85d6376aa1fd7422fa9f36899f99 Mon Sep 17 00:00:00 2001 From: Simon Binder Date: Wed, 2 Sep 2026 10:39:34 +0200 Subject: [PATCH 07/15] AI feedback --- .../evaluator/parameter_evaluator.ts | 8 ++++-- .../src/sync_plan/evaluator/result_set.ts | 25 +++++++++++-------- 2 files changed, 20 insertions(+), 13 deletions(-) diff --git a/packages/sync-rules/src/sync_plan/evaluator/parameter_evaluator.ts b/packages/sync-rules/src/sync_plan/evaluator/parameter_evaluator.ts index a1a9175f3..617312702 100644 --- a/packages/sync-rules/src/sync_plan/evaluator/parameter_evaluator.ts +++ b/packages/sync-rules/src/sync_plan/evaluator/parameter_evaluator.ts @@ -172,7 +172,11 @@ export class RequestParameterEvaluators { } } - return this.#readParameters()!; + const params = this.#readParameters(); + if (params == null) { + throw new Error('internal error: Should have been able to resolve instantiation after instantiating stages.'); + } + return params; } catch (e) { if (e === uninstantiableException) return []; @@ -228,7 +232,7 @@ export class RequestParameterEvaluators { for (const value of constraint.values) { if (value instanceof RequestParameterValue) { const evaluated = value.requireResolved(); - if (knownValue !== undefined && evaluated != knownValue) { + if (knownValue !== undefined && evaluated !== knownValue) { throw uninstantiableException; } diff --git a/packages/sync-rules/src/sync_plan/evaluator/result_set.ts b/packages/sync-rules/src/sync_plan/evaluator/result_set.ts index 6c41cddf0..264bf8b76 100644 --- a/packages/sync-rules/src/sync_plan/evaluator/result_set.ts +++ b/packages/sync-rules/src/sync_plan/evaluator/result_set.ts @@ -32,7 +32,7 @@ export class ResultSet { for (const row of this.#rows) { // We can shallow-clone rows, inner items are frozen once added into the result set. - rs.#rows.push([...row]); + rs.#rows.push(row.slice()); } return rs; } @@ -108,22 +108,25 @@ export class ResultSet { * If a fixed value is passed, this also removes rows where any of the given columns has a different value. */ formIntersection(columns: ResultSetColumn[], fixedValue?: SqliteParameterValue) { - row: for (let i = 0; i < this.#rows.length; i++) { - const row = this.#rows[i]; + const keptRows: ResultSetRow[] = []; + + row: for (const row of this.#rows) { let requiredValue = fixedValue; for (const column of columns) { const evaluated = lookupInRow(row, column); - if (requiredValue !== undefined && evaluated != requiredValue) { - // This row needs to be removed! - this.#rows.splice(i, 1); - i--; + if (requiredValue !== undefined && evaluated !== requiredValue) { + // Intersection doesn't match, skip this row. continue row; } requiredValue = evaluated; } + + keptRows.push(row); } + + this.#rows = keptRows; } #multiplyAtRow(resultSetIndex: number, rowIndex: number, rows: SqliteParameterValue[][]) { @@ -132,7 +135,7 @@ export class ResultSet { row[resultSetIndex] = Object.freeze(rows[0]); for (let j = 1; j < rows.length; j++) { - const copy = Array.from(row); + const copy = row.slice(); copy[resultSetIndex] = Object.freeze(rows[j]); this.#rows.push(copy); } @@ -152,11 +155,11 @@ export class ResultSet { const existingGroup = foundValues.get(value); if (existingGroup != null) { - yield { index: i, group: existingGroup, first: false }; + yield { group: existingGroup, first: false }; } else { const group = generateGroup([value]); foundValues.set(value, group); - yield { index: i, group, first: true }; + yield { group, first: true }; } } } else { @@ -172,7 +175,7 @@ export class ResultSet { return generateGroup(values); }); - yield { index: i, group, first: isFirst }; + yield { group, first: isFirst }; } } } From d5500feb419cc628a3f6af8a7a8f813fd1ba573d Mon Sep 17 00:00:00 2001 From: Simon Binder Date: Wed, 2 Sep 2026 11:10:44 +0200 Subject: [PATCH 08/15] Evaluate intersection on table-valued functions eagerly --- .../src/sync_plan/evaluator/parameter_evaluator.ts | 4 +++- .../test/src/sync_plan/evaluator/evaluator.test.ts | 8 ++++---- 2 files changed, 7 insertions(+), 5 deletions(-) diff --git a/packages/sync-rules/src/sync_plan/evaluator/parameter_evaluator.ts b/packages/sync-rules/src/sync_plan/evaluator/parameter_evaluator.ts index 617312702..22302428d 100644 --- a/packages/sync-rules/src/sync_plan/evaluator/parameter_evaluator.ts +++ b/packages/sync-rules/src/sync_plan/evaluator/parameter_evaluator.ts @@ -129,7 +129,7 @@ export class RequestParameterEvaluators { for (const instantiation of stage.inputParameters()) { if (instantiation instanceof RequestParameterValue) { instantiation.resolveWith(input); - } else { + } else if (instantiation.lookup instanceof ParameterIndexExpandingLookup) { needsParameterLookups = true; } } @@ -225,6 +225,8 @@ export class RequestParameterEvaluators { } #applyIntersectionConstraint(constraint: RequiredIntersection) { + if (constraint.wasApplied) return; + // If any parameter of the intersection is a scalar value derived from a request, that value. let knownValue: SqliteParameterValue | undefined; const intersection: ResultSetColumn[] = []; diff --git a/packages/sync-rules/test/src/sync_plan/evaluator/evaluator.test.ts b/packages/sync-rules/test/src/sync_plan/evaluator/evaluator.test.ts index 8cc13553b..1b70011ed 100644 --- a/packages/sync-rules/test/src/sync_plan/evaluator/evaluator.test.ts +++ b/packages/sync-rules/test/src/sync_plan/evaluator/evaluator.test.ts @@ -558,10 +558,10 @@ config: streams: stream: auto_subscribe: true - query: SELECT * FROM issues WHERE a = auth.parameter('x') AND a = auth.parameter('y') + query: SELECT * FROM issues WHERE a = auth.parameter('x') AND a IN auth.parameter('y') `); - function queryWith(x: string, y: string) { + function queryWith(x: string, y: string[]) { const { querier, errors } = desc.getBucketParameterQuerier({ globalParameters: requestParameters({ sub: 'user', x, y }), hasDefaultStreams: true, @@ -572,8 +572,8 @@ streams: return querier.staticBuckets.map((e) => e.bucket); } - expect(queryWith('p1', 'p2')).toStrictEqual([]); - expect(queryWith('p1', 'p1')).toStrictEqual(['stream|0["p1"]']); + expect(queryWith('p1', ['p2'])).toStrictEqual([]); + expect(queryWith('p1', ['p1', 'p2'])).toStrictEqual(['stream|0["p1"]']); }); syncTest('parameter lookups', async ({ sync }) => { From 849a419d50c8b457a89e5bf8d8350372ec4f907a Mon Sep 17 00:00:00 2001 From: Simon Binder Date: Tue, 8 Sep 2026 15:22:08 +0200 Subject: [PATCH 09/15] Build intersection row-by-row --- .../evaluator/parameter_evaluator.ts | 100 ++++------ .../src/sync_plan/evaluator/result_set.ts | 183 ++++++++++++++---- .../sync_plan/evaluator/result_set.test.ts | 116 +++++++++-- 3 files changed, 274 insertions(+), 125 deletions(-) diff --git a/packages/sync-rules/src/sync_plan/evaluator/parameter_evaluator.ts b/packages/sync-rules/src/sync_plan/evaluator/parameter_evaluator.ts index 22302428d..83eb6395f 100644 --- a/packages/sync-rules/src/sync_plan/evaluator/parameter_evaluator.ts +++ b/packages/sync-rules/src/sync_plan/evaluator/parameter_evaluator.ts @@ -13,7 +13,7 @@ import { MapSourceVisitor, visitExpr } from '../expression_visitor.js'; import * as plan from '../plan.js'; import { StreamInput } from './bucket_source.js'; import { PreparedParameterIndexLookupCreator } from './parameter_index_lookup_creator.js'; -import { AsyncJoinLookup, ResultSet, ResultSetColumn, ResultSetElement } from './result_set.js'; +import { AsyncJoinLookup, IntersectionConstraint, ResultSet, ResultSetColumn, ResultSetElement } from './result_set.js'; /** * Finds bucket parameters for a given request or subscription. @@ -64,7 +64,10 @@ export class RequestParameterEvaluators { * Pending parameter values, or their cached outputs. */ private readonly parameterValues: PreparedParameterValue[], - + /** + * Intersection constraints to take into consideration when adding rows to the result set. + */ + private readonly intersectionConstraints: PendingIntersectionConstraint[], /** * The materialized result set containing lookup values. {@link parameterValues} are read from this result set as a * final step. @@ -97,7 +100,13 @@ export class RequestParameterEvaluators { const copiedStages = this.lookupStages.map((s) => s.clone(cloneParameter)); const outputValues = this.parameterValues.map(cloneParameter); - return new RequestParameterEvaluators(this.stream, copiedStages, outputValues, this.resultSet.clone()); + return new RequestParameterEvaluators( + this.stream, + copiedStages, + outputValues, + this.intersectionConstraints.slice(), + this.resultSet.clone() + ); } /** @@ -110,10 +119,12 @@ export class RequestParameterEvaluators { */ partiallyInstantiate(input: PartialInstantiationInput): SqliteParameterValue[][] | undefined { try { + this.resultSet.addIntersectionConstraints( + this.intersectionConstraints.map((c) => this.#createIntersectionResult(c, input)) + ); + // At this point, we can resolve table-valued lookups and parameter values based only on request data. for (const stage of this.lookupStages) { - let needsParameterLookups = false; - for (const element of stage.lookups) { if (element instanceof TableValuedExpandingLookup) { const outputs = element.read(input.request); @@ -121,22 +132,12 @@ export class RequestParameterEvaluators { this.resultSet.multiply(element.resultSetIndex, outputs); this.#checkInstantiable(); - } else { - needsParameterLookups = true; } } for (const instantiation of stage.inputParameters()) { if (instantiation instanceof RequestParameterValue) { instantiation.resolveWith(input); - } else if (instantiation.lookup instanceof ParameterIndexExpandingLookup) { - needsParameterLookups = true; - } - } - - if (!needsParameterLookups) { - for (const intersection of stage.intersections) { - this.#applyIntersectionConstraint(intersection); } } } @@ -160,16 +161,12 @@ export class RequestParameterEvaluators { */ async instantiate(input: InstantiationInput): Promise { try { - for (const { lookups, intersections } of this.lookupStages) { + for (const { lookups } of this.lookupStages) { for (const lookup of lookups) { if (lookup instanceof ParameterIndexExpandingLookup) { await this.#instantiateLookup(lookup, input); } } - - for (const intersection of intersections) { - this.#applyIntersectionConstraint(intersection); - } } const params = this.#readParameters(); @@ -189,11 +186,7 @@ export class RequestParameterEvaluators { } #readParameters(): SqliteParameterValue[][] | undefined { - for (const { intersections, lookups } of this.lookupStages) { - for (const intersection of intersections) { - if (!intersection.wasApplied) return undefined; - } - + for (const { lookups } of this.lookupStages) { for (const element of lookups) { if (!element.wasResolved) return undefined; } @@ -224,29 +217,28 @@ export class RequestParameterEvaluators { }); } - #applyIntersectionConstraint(constraint: RequiredIntersection) { - if (constraint.wasApplied) return; - + #createIntersectionResult( + constraint: PendingIntersectionConstraint, + input: PartialInstantiationInput + ): IntersectionConstraint { // If any parameter of the intersection is a scalar value derived from a request, that value. - let knownValue: SqliteParameterValue | undefined; + let fixedValue: SqliteParameterValue | undefined; const intersection: ResultSetColumn[] = []; - for (const value of constraint.values) { + for (const value of constraint.inputs) { if (value instanceof RequestParameterValue) { - const evaluated = value.requireResolved(); - if (knownValue !== undefined && evaluated !== knownValue) { + const evaluated = value.resolveWith(input); + if (fixedValue !== undefined && evaluated !== fixedValue) { throw uninstantiableException; } - knownValue = evaluated; + fixedValue = evaluated; } else { intersection.push(value); } } - this.resultSet.formIntersection(intersection, knownValue); - constraint.wasApplied = true; - this.#checkInstantiable(); + return { fixedValue, columns: intersection }; } async #instantiateLookup(lookup: ParameterIndexExpandingLookup, input: InstantiationInput) { @@ -312,6 +304,7 @@ export class RequestParameterEvaluators { const mappedStages: LookupStage[] = []; let amountOfLookups = 0; const lookupToStage = new Map(); + const intersections: PendingIntersectionConstraint[] = []; function mapParameterValue(value: plan.ParameterValue): PreparedParameterValue { if (value.type == 'request') { @@ -328,14 +321,7 @@ export class RequestParameterEvaluators { return new LookupParameterValue(lookup, value.resultIndex); } else { const intersectionInputs = mapParameterValues(value.values); - - if (mappedStages.length > 0) { - mappedStages[mappedStages.length - 1].intersections.push({ values: intersectionInputs, wasApplied: false }); - } else { - // Intersection in first stage, e.g. for request parameters. Add a stage just for this. - const stage = new LookupStage([], [{ values: intersectionInputs, wasApplied: false }]); - mappedStages.push(stage); - } + intersections.push({ inputs: intersectionInputs }); // Non-intersecting rows will be pruned from the result set or, for scalar parameter values, mark the querier // as uninstantiable. So, we can replace the intersection value with any inner value. @@ -348,7 +334,7 @@ export class RequestParameterEvaluators { } for (const stage of lookupStages) { - const mappedStage = new LookupStage([], []); + const mappedStage = new LookupStage([]); for (const lookup of stage) { let resolved: PreparedExpandingLookup; @@ -392,7 +378,7 @@ export class RequestParameterEvaluators { } const rs = new ResultSet(amountOfLookups); - return new RequestParameterEvaluators(stream, mappedStages, mapParameterValues(values), rs); + return new RequestParameterEvaluators(stream, mappedStages, mapParameterValues(values), intersections, rs); } } @@ -455,20 +441,10 @@ class LookupStage { /** * Lookups that only have dependencies on prior stages. */ - readonly lookups: PreparedExpandingLookup[], - /** - * A list of constraints enforcing that specific columns must have equal values. - * - * These constraints are evaluated after lookups, and may reference lookups in this stage. - */ - readonly intersections: RequiredIntersection[] + readonly lookups: PreparedExpandingLookup[] ) {} *inputParameters() { - for (const intersection of this.intersections) { - yield* intersection.values; - } - for (const lookup of this.lookups) { if (lookup instanceof ParameterIndexExpandingLookup) { yield* lookup.instantiation; @@ -477,16 +453,12 @@ class LookupStage { } clone(cloneParameter: CloneParameter): LookupStage { - return new LookupStage( - this.lookups.map((l) => l.clone(cloneParameter)), - this.intersections.map(({ values, wasApplied }) => ({ values: values.map(cloneParameter), wasApplied })) - ); + return new LookupStage(this.lookups.map((l) => l.clone(cloneParameter))); } } -interface RequiredIntersection { - values: PreparedParameterValue[]; - wasApplied: boolean; +interface PendingIntersectionConstraint { + inputs: PreparedParameterValue[]; } /** diff --git a/packages/sync-rules/src/sync_plan/evaluator/result_set.ts b/packages/sync-rules/src/sync_plan/evaluator/result_set.ts index 264bf8b76..c2fab55df 100644 --- a/packages/sync-rules/src/sync_plan/evaluator/result_set.ts +++ b/packages/sync-rules/src/sync_plan/evaluator/result_set.ts @@ -9,14 +9,24 @@ import { SqliteParameterValue } from '../../types.js'; * by reading columns in each row. */ export class ResultSet { - #totalLookups: number; + #containsLookup: boolean[]; #rows: ResultSetRow[]; + /** + * Intersection constraints to respect when adding new rows, keyed by result sets affected by them. + * + * Invariant: For all lookups that have already been added, no row contradicts any intersection constraint. In other + * words, we only have to check _one_ existing column of the intersection when processing new lookups to add. + */ + #intersections: (IntersectionConstraint[] | undefined)[]; + /** * @param totalLookups - The total amount of lookups that will be joined to this result set. */ constructor(totalLookups: number) { - this.#totalLookups = totalLookups; + this.#containsLookup = new Array(totalLookups).fill(false); + this.#intersections = new Array(totalLookups).fill(undefined); + const initialRow = new Array(totalLookups); initialRow.fill(undefined); this.#rows = [initialRow]; @@ -27,16 +37,38 @@ export class ResultSet { } clone(): ResultSet { - const rs = new ResultSet(this.#totalLookups); + const rs = new ResultSet(this.#containsLookup.length); + rs.#containsLookup = this.#containsLookup.slice(); + rs.#intersections = this.#intersections.slice(); rs.#rows.splice(0, 1); // Remove the initial unit row for (const row of this.#rows) { // We can shallow-clone rows, inner items are frozen once added into the result set. rs.#rows.push(row.slice()); } + return rs; } + addIntersectionConstraints(constraints: Iterable) { + // To make it easy to uphold the invariant that all rows must satisfy the constraint, we only allow adding + // intersection constraints to the initial result set. + if (this.length != 1 || this.#rows[0].some((s) => s !== undefined)) { + throw new Error('Can only add intersection constraints to unit result set'); + } + + for (const constraint of constraints) { + for (const { lookup } of constraint.columns) { + const tracked = this.#intersections[lookup.resultSetIndex]; + if (tracked != null && !tracked.includes(constraint)) { + tracked.push(constraint); + } else { + this.#intersections[lookup.resultSetIndex] = [constraint]; + } + } + } + } + /** * Extracts unique values by looking values for each column in this result set. */ @@ -54,11 +86,16 @@ export class ResultSet { multiply(resultSetIndex: number, rows: SqliteParameterValue[][]) { if (rows.length === 0) { this.#rows = []; + return; } + using add = this.#prepareAddingResultSet(resultSetIndex); const originalLength = this.#rows.length; + for (let i = 0; i < originalLength; i++) { - this.#multiplyAtRow(resultSetIndex, i, rows); + if (!this.#multiplyAtRow(resultSetIndex, add.filter, i, rows)) { + add.deletedRows.push(i); + } } } @@ -82,63 +119,120 @@ export class ResultSet { await performLookup(uniqueLookups); - const deletedRows: number[] = []; + using add = this.#prepareAddingResultSet(resultSetIndex); + const originalLength = this.#rows.length; for (let i = 0; i < originalLength; i++) { const lookup = lookupsByRow[i]; - if (lookup.foundRows.length > 0) { - this.#multiplyAtRow(resultSetIndex, i, lookup.foundRows); - } else { + if (lookup.foundRows.length === 0 || !this.#multiplyAtRow(resultSetIndex, add.filter, i, lookup.foundRows)) { // The row has no matching join partner, so remove it. We can't split it immediately because #multiplyAtRow is // still iterating through rows. - deletedRows.push(i); + add.deletedRows.push(i); } } - - let offset = 0; - for (const toDelete of deletedRows) { - this.#rows.splice(toDelete - offset, 1); - offset++; - } } - /** - * Removes rows where the given columns have different values. - * - * If a fixed value is passed, this also removes rows where any of the given columns has a different value. - */ - formIntersection(columns: ResultSetColumn[], fixedValue?: SqliteParameterValue) { - const keptRows: ResultSetRow[] = []; + #intersectionFilter(intersection: IntersectionConstraint, addedResultSetIndex: number): IntersectionFilter { + const affectedColumnsInAddedResultSet = intersection.columns.filter( + ({ lookup }) => lookup.resultSetIndex === addedResultSetIndex + ); + + if (intersection.fixedValue !== undefined) { + // All rows already in the result set satisfy the intersection and must match the fixed value in relevant columns. + // So when checking a new row, we just need to check columns there. + return function (_existingRow: ResultSetRow, added: SqliteParameterValue[]): boolean { + for (const { outputIndex } of affectedColumnsInAddedResultSet) { + if (added[outputIndex] !== intersection.fixedValue) return false; + } - row: for (const row of this.#rows) { - let requiredValue = fixedValue; + return true; + }; + } else { + const anyExistingColumn = intersection.columns.find(({ lookup }) => this.#containsLookup[lookup.resultSetIndex]); + + return function (existingRow: ResultSetRow, added: SqliteParameterValue[]): boolean { + let referenceValue: SqliteParameterValue | undefined; - for (const column of columns) { - const evaluated = lookupInRow(row, column); - if (requiredValue !== undefined && evaluated !== requiredValue) { - // Intersection doesn't match, skip this row. - continue row; + if (anyExistingColumn != null) { + referenceValue = lookupInRow(existingRow, anyExistingColumn); } - requiredValue = evaluated; - } + for (const { outputIndex } of affectedColumnsInAddedResultSet) { + const value = added[outputIndex]; - keptRows.push(row); + if (referenceValue !== undefined && value !== referenceValue) return false; + referenceValue = value; + } + + return true; + }; } + } - this.#rows = keptRows; + #prepareAddingResultSet(addedResultSetIndex: number) { + if (this.#containsLookup[addedResultSetIndex]) { + throw new Error(`Already added results for ${addedResultSetIndex}`); + } + + const filters = this.#intersections[addedResultSetIndex]?.map((intersection) => + this.#intersectionFilter(intersection, addedResultSetIndex) + ); + + const deletedRows: number[] = []; + const filter = (existingRow: ResultSetRow, added: SqliteParameterValue[]): boolean => { + if (filters == null) return true; + + return filters.every((f) => f(existingRow, added)); + }; + + return { + deletedRows, + filter, + [Symbol.dispose]: () => { + let offset = 0; + for (const toDelete of deletedRows) { + this.#rows.splice(toDelete - offset, 1); + offset++; + } + + this.#containsLookup[addedResultSetIndex] = true; + } + }; } - #multiplyAtRow(resultSetIndex: number, rowIndex: number, rows: SqliteParameterValue[][]) { - // Add first element of product to existing row, remaining as new rows. - const row = this.#rows[rowIndex]; - row[resultSetIndex] = Object.freeze(rows[0]); + /** + * Adds the cartesian product of an existing row and a new result set. + * + * Returns false if the original row needs to be removed because no row was added (e.g. because an intersection filter + * doesn't match). + */ + #multiplyAtRow( + resultSetIndex: number, + filter: IntersectionFilter, + rowIndex: number, + rows: SqliteParameterValue[][] + ): boolean { + const originalRow = this.#rows[rowIndex]; + + let isFirst = true; + for (const row of rows) { + if (!filter(originalRow, row)) { + continue; + } + + if (isFirst) { + isFirst = false; - for (let j = 1; j < rows.length; j++) { - const copy = row.slice(); - copy[resultSetIndex] = Object.freeze(rows[j]); - this.#rows.push(copy); + // Add first element of product to existing row, remaining as new rows. + originalRow[resultSetIndex] = Object.freeze(row); + } else { + const copy = originalRow.slice(); + copy[resultSetIndex] = Object.freeze(row); + this.#rows.push(copy); + } } + + return !isFirst; } *#groupBy(columns: ResultSetColumn[], generateGroup: (values: SqliteParameterValue[]) => T) { @@ -195,6 +289,13 @@ export interface AsyncJoinLookup { foundRows: SqliteParameterValue[][]; } +export interface IntersectionConstraint { + readonly columns: ResultSetColumn[]; + readonly fixedValue?: SqliteParameterValue; +} + +type IntersectionFilter = (existingRow: ResultSetRow, added: SqliteParameterValue[]) => boolean; + /** * A row in a result set. * diff --git a/packages/sync-rules/test/src/sync_plan/evaluator/result_set.test.ts b/packages/sync-rules/test/src/sync_plan/evaluator/result_set.test.ts index 27490cf55..ffe2e3c94 100644 --- a/packages/sync-rules/test/src/sync_plan/evaluator/result_set.test.ts +++ b/packages/sync-rules/test/src/sync_plan/evaluator/result_set.test.ts @@ -71,33 +71,53 @@ describe('ResultSet', () => { ['b', 1] ]); }); - }); - describe('formIntersection', () => { - const col0 = { lookup: element(0), outputIndex: 0 }; - const col1 = { lookup: element(1), outputIndex: 0 }; + describe('respects intersection', () => { + test('with constant', () => { + const rs = new ResultSet(1); - test('with a fixed value, removes rows where any column differs from it', () => { - const rs = new ResultSet(2); - rs.multiply(0, [['a'], ['b']]); - rs.multiply(1, [['a'], ['b']]); + rs.addIntersectionConstraints([{ fixedValue: 'test', columns: [{ lookup: element(0), outputIndex: 0 }] }]); + rs.multiply(0, [['test'], ['otherValue']]); + expect(rs.length).toStrictEqual(1); + }); - rs.formIntersection([col0, col1], 'a'); + test('with column in same row', () => { + const rs = new ResultSet(1); - expect([...rs.projectUnique([col0, col1])]).toStrictEqual([['a', 'a']]); - }); + rs.addIntersectionConstraints([ + { + columns: [ + { lookup: element(0), outputIndex: 0 }, + { lookup: element(0), outputIndex: 1 } + ] + } + ]); + rs.multiply(0, [ + ['a', 'a'], + ['a', 'b'], + ['b', 'a'], + ['b', 'b'] + ]); + expect(rs.length).toStrictEqual(2); + }); - test('without a fixed value, removes rows where the columns differ from each other', () => { - const rs = new ResultSet(2); - rs.multiply(0, [['a'], ['b']]); - rs.multiply(1, [['a'], ['b']]); + test('with existing column', () => { + const rs = new ResultSet(2); - rs.formIntersection([col0, col1]); + rs.addIntersectionConstraints([ + { + columns: [ + { lookup: element(0), outputIndex: 0 }, + { lookup: element(1), outputIndex: 0 } + ] + } + ]); - expect([...rs.projectUnique([col0, col1])]).toStrictEqual([ - ['a', 'a'], - ['b', 'b'] - ]); + rs.multiply(0, [['test0'], ['test1']]); + expect(rs.length).toStrictEqual(2); + rs.multiply(1, [['test0'], ['unrelated']]); + expect(rs.length).toStrictEqual(1); + }); }); }); @@ -215,6 +235,62 @@ describe('ResultSet', () => { ['b', 'x'] ]); }); + + describe('respects intersection', () => { + test('with constant', async () => { + const rs = new ResultSet(1); + + rs.addIntersectionConstraints([{ fixedValue: 'test', columns: [{ lookup: element(0), outputIndex: 0 }] }]); + await rs.joinAsync([], 0, async (lookups) => { + lookups[0].foundRows.push(['test'], ['otherValue']); + }); + + expect(rs.length).toStrictEqual(1); + }); + + test('with column in same row', async () => { + const rs = new ResultSet(1); + + rs.addIntersectionConstraints([ + { + columns: [ + { lookup: element(0), outputIndex: 0 }, + { lookup: element(0), outputIndex: 1 } + ] + } + ]); + await rs.joinAsync([], 0, async (lookups) => { + lookups[0].foundRows.push(['a', 'a'], ['a', 'b'], ['b', 'a'], ['b', 'b']); + }); + + expect(rs.length).toStrictEqual(2); + }); + + test('with existing column', async () => { + const rs = new ResultSet(2); + + rs.addIntersectionConstraints([ + { + columns: [ + { lookup: element(0), outputIndex: 0 }, + { lookup: element(1), outputIndex: 0 } + ] + } + ]); + + rs.multiply(0, [['test0'], ['test1']]); + expect(rs.length).toStrictEqual(2); + + // Both rows share the same (empty) join key, so the lookup is only performed once even though the + // intersection filter must still be checked per-row against the existing column's value. + await rs.joinAsync([], 1, async (lookups) => { + expect(lookups.length).toStrictEqual(1); + lookups[0].foundRows.push(['test0'], ['unrelated']); + }); + + expect(rs.length).toStrictEqual(1); + }); + }); }); }); From b5cf4eb1b399fb0deb79430afead46161c82b5ab Mon Sep 17 00:00:00 2001 From: Simon Binder Date: Tue, 8 Sep 2026 15:37:08 +0200 Subject: [PATCH 10/15] AI feedback --- .../src/sync_plan/evaluator/parameter_evaluator.ts | 2 +- packages/sync-rules/src/sync_plan/evaluator/result_set.ts | 6 ++++-- 2 files changed, 5 insertions(+), 3 deletions(-) diff --git a/packages/sync-rules/src/sync_plan/evaluator/parameter_evaluator.ts b/packages/sync-rules/src/sync_plan/evaluator/parameter_evaluator.ts index 83eb6395f..5ee59fb26 100644 --- a/packages/sync-rules/src/sync_plan/evaluator/parameter_evaluator.ts +++ b/packages/sync-rules/src/sync_plan/evaluator/parameter_evaluator.ts @@ -104,7 +104,7 @@ export class RequestParameterEvaluators { this.stream, copiedStages, outputValues, - this.intersectionConstraints.slice(), + this.intersectionConstraints.map(({ inputs }) => ({ inputs: inputs.map(cloneParameter) })), this.resultSet.clone() ); } diff --git a/packages/sync-rules/src/sync_plan/evaluator/result_set.ts b/packages/sync-rules/src/sync_plan/evaluator/result_set.ts index c2fab55df..61d06501f 100644 --- a/packages/sync-rules/src/sync_plan/evaluator/result_set.ts +++ b/packages/sync-rules/src/sync_plan/evaluator/result_set.ts @@ -60,8 +60,10 @@ export class ResultSet { for (const constraint of constraints) { for (const { lookup } of constraint.columns) { const tracked = this.#intersections[lookup.resultSetIndex]; - if (tracked != null && !tracked.includes(constraint)) { - tracked.push(constraint); + if (tracked != null) { + if (!tracked.includes(constraint)) { + tracked.push(constraint); + } } else { this.#intersections[lookup.resultSetIndex] = [constraint]; } From a5736519e85e4d3ccb57f074eed800b14bc4564b Mon Sep 17 00:00:00 2001 From: Simon Binder Date: Thu, 10 Sep 2026 18:25:04 +0200 Subject: [PATCH 11/15] Avoid quadratic result set cleanup --- .../src/sync_plan/evaluator/result_set.ts | 20 +++++++++++++++---- 1 file changed, 16 insertions(+), 4 deletions(-) diff --git a/packages/sync-rules/src/sync_plan/evaluator/result_set.ts b/packages/sync-rules/src/sync_plan/evaluator/result_set.ts index 61d06501f..32a20ba64 100644 --- a/packages/sync-rules/src/sync_plan/evaluator/result_set.ts +++ b/packages/sync-rules/src/sync_plan/evaluator/result_set.ts @@ -191,10 +191,22 @@ export class ResultSet { deletedRows, filter, [Symbol.dispose]: () => { - let offset = 0; - for (const toDelete of deletedRows) { - this.#rows.splice(toDelete - offset, 1); - offset++; + if (deletedRows.length > 0) { + // Indices to delete are sorted, so we can do a single pass to the array to remove them. We fill deleted slots + // by moving following rows into the hole. + let deletedRowsIndex = 0; + let writeIndex = deletedRows[0]; + + for (let scanIndex = deletedRows[0]; scanIndex < this.#rows.length; scanIndex++) { + if (deletedRowsIndex < deletedRows.length && scanIndex === deletedRows[deletedRowsIndex]) { + // Skip this row to delete it, the next row will be copied to this position. + deletedRowsIndex++; + continue; + } + this.#rows[writeIndex++] = this.#rows[scanIndex]; + } + + this.#rows.length = writeIndex; } this.#containsLookup[addedResultSetIndex] = true; From 84a752f377e7ec85f348affd4957c2279ff8892a Mon Sep 17 00:00:00 2001 From: Simon Binder Date: Mon, 14 Sep 2026 09:09:02 +0200 Subject: [PATCH 12/15] Make more efficient because we don't care about order --- .../src/sync_plan/evaluator/result_set.ts | 21 ++++++++----------- .../sync_plan/evaluator/result_set.test.ts | 4 ++-- 2 files changed, 11 insertions(+), 14 deletions(-) diff --git a/packages/sync-rules/src/sync_plan/evaluator/result_set.ts b/packages/sync-rules/src/sync_plan/evaluator/result_set.ts index 32a20ba64..2d05f6490 100644 --- a/packages/sync-rules/src/sync_plan/evaluator/result_set.ts +++ b/packages/sync-rules/src/sync_plan/evaluator/result_set.ts @@ -192,21 +192,18 @@ export class ResultSet { filter, [Symbol.dispose]: () => { if (deletedRows.length > 0) { - // Indices to delete are sorted, so we can do a single pass to the array to remove them. We fill deleted slots - // by moving following rows into the hole. - let deletedRowsIndex = 0; - let writeIndex = deletedRows[0]; - - for (let scanIndex = deletedRows[0]; scanIndex < this.#rows.length; scanIndex++) { - if (deletedRowsIndex < deletedRows.length && scanIndex === deletedRows[deletedRowsIndex]) { - // Skip this row to delete it, the next row will be copied to this position. - deletedRowsIndex++; - continue; + // This is a set, so the order of rows doesn't matter. Delete indices by moving them to the end, then + // truncating. Note that deletedRows are sorted, this processes them from the end. + let end = this.#rows.length; + for (let i = deletedRows.length - 1; i >= 0; i--) { + const deleted = deletedRows[i]; + end--; + if (deleted !== end) { + this.#rows[deleted] = this.#rows[end]; } - this.#rows[writeIndex++] = this.#rows[scanIndex]; } - this.#rows.length = writeIndex; + this.#rows.length = end; } this.#containsLookup[addedResultSetIndex] = true; diff --git a/packages/sync-rules/test/src/sync_plan/evaluator/result_set.test.ts b/packages/sync-rules/test/src/sync_plan/evaluator/result_set.test.ts index ffe2e3c94..0a9c0bfec 100644 --- a/packages/sync-rules/test/src/sync_plan/evaluator/result_set.test.ts +++ b/packages/sync-rules/test/src/sync_plan/evaluator/result_set.test.ts @@ -172,8 +172,8 @@ describe('ResultSet', () => { expect(rs.length).toStrictEqual(4); expect([...rs.projectUnique([col0, col1])]).toStrictEqual([ ['a', 10], - ['c', 30], - ['c', 31] + ['c', 31], + ['c', 30] ]); }); From c9bf009cf7f0f28cb1b9496ff8cb64be6f69197a Mon Sep 17 00:00:00 2001 From: Simon Binder Date: Mon, 14 Sep 2026 11:06:00 +0200 Subject: [PATCH 13/15] Optimizations --- .../evaluator/parameter_evaluator.ts | 121 ++++++++++++++++-- .../src/sync_plan/evaluator/evaluator.test.ts | 118 ++++++++++++++++- .../test/src/sync_plan/evaluator/utils.ts | 8 +- 3 files changed, 229 insertions(+), 18 deletions(-) diff --git a/packages/sync-rules/src/sync_plan/evaluator/parameter_evaluator.ts b/packages/sync-rules/src/sync_plan/evaluator/parameter_evaluator.ts index 5ee59fb26..af0b2ced5 100644 --- a/packages/sync-rules/src/sync_plan/evaluator/parameter_evaluator.ts +++ b/packages/sync-rules/src/sync_plan/evaluator/parameter_evaluator.ts @@ -119,19 +119,102 @@ export class RequestParameterEvaluators { */ partiallyInstantiate(input: PartialInstantiationInput): SqliteParameterValue[][] | undefined { try { - this.resultSet.addIntersectionConstraints( - this.intersectionConstraints.map((c) => this.#createIntersectionResult(c, input)) - ); + const mappedIntersectionConstraints = new Set(); + + // Optimization: For request-parameter based lookups with a single column (the most common kind), use a set to + // immediately filter out rows that don't intersect with another request parameter. ResultSet.multiply does the + // same thing, but that requires going through the combination of all rows whereas the de-duplication here is + // linear. + interface RequestParameterIntersection { + constraint: IntersectionConstraint; + // If the intersection consists entirely of request parameters, and all mentioned request parameters are in + // exactly one intersection, we don't need to register the constraint on the result set as it's covered by the + // optimization. + replacesConstraint: boolean; + amountOfRequestParameters: number; + // Count in how many parameters a value was seen, we filter out values not present in all parameters. + rows: Map; + materializedRows: SqliteParameterValue[][]; + } + + const parameterToIntersection = new Map(); + + for (const constraint of this.intersectionConstraints) { + const mapped = this.#createIntersectionResult(constraint, input); + mappedIntersectionConstraints.add(mapped); + + const parameterIntersection: RequestParameterIntersection = { + constraint: mapped, + replacesConstraint: true, + amountOfRequestParameters: 0, + rows: new Map(), + materializedRows: [] + }; + + for (const column of constraint.inputs) { + if ( + column instanceof LookupParameterValue && + column.lookup instanceof TableValuedExpandingLookup && + column.lookup.columnCount == 1 + ) { + parameterIntersection.amountOfRequestParameters++; + + const existing = parameterToIntersection.get(column.lookup); + + if (existing != null) { + existing.push(parameterIntersection); + + // This parameter is part of multiple intersections. + for (const intersection of existing) { + intersection.replacesConstraint = false; + } + } else { + parameterToIntersection.set(column.lookup, [parameterIntersection]); + } + } + } + + if (parameterIntersection.amountOfRequestParameters < constraint.inputs.length) { + // This intersection doesn't entirely consist of request parameters. + parameterIntersection.replacesConstraint = false; + } + } + + for (const intersections of parameterToIntersection.values()) { + for (const intersection of intersections) { + if (intersection.replacesConstraint) { + mappedIntersectionConstraints.delete(intersection.constraint); + } + } + } + + this.resultSet.addIntersectionConstraints(mappedIntersectionConstraints); // At this point, we can resolve table-valued lookups and parameter values based only on request data. for (const stage of this.lookupStages) { for (const element of stage.lookups) { if (element instanceof TableValuedExpandingLookup) { const outputs = element.read(input.request); - element.wasResolved = true; - this.resultSet.multiply(element.resultSetIndex, outputs); - - this.#checkInstantiable(); + const intersections = parameterToIntersection.get(element); + if (intersections == null || intersections.length == 0) { + this.resultSet.multiply(element.resultSetIndex, outputs); + this.#checkInstantiable(); + } else { + // For this parameter to be part of the intersection optimization, it must return exactly one column. + for (const [value] of outputs) { + for (const intersection of intersections) { + const fixed = intersection.constraint.fixedValue; + if (fixed != null && value != fixed) continue; + + const matchedSources = (intersection.rows.get(value) ?? 0) + 1; + + intersection.rows.set(value, matchedSources); + if (matchedSources == intersection.amountOfRequestParameters) { + intersection.materializedRows.push([value]); + } + } + } + } } } @@ -142,6 +225,16 @@ export class RequestParameterEvaluators { } } + for (const [parameter, intersections] of parameterToIntersection.entries()) { + // Find the smallest intersection constraining this parameter, then instantiate the parameter to that. + const smallestIntersection = intersections.reduce((acc, intersection) => + intersection.materializedRows.length < acc.materializedRows.length ? intersection : acc + ); + + this.resultSet.multiply(parameter.resultSetIndex, smallestIntersection.materializedRows); + this.#checkInstantiable(); + } + for (const parameter of this.parameterValues) { if (parameter instanceof RequestParameterValue) parameter.resolveWith(input); } @@ -188,7 +281,7 @@ export class RequestParameterEvaluators { #readParameters(): SqliteParameterValue[][] | undefined { for (const { lookups } of this.lookupStages) { for (const element of lookups) { - if (!element.wasResolved) return undefined; + if (element instanceof ParameterIndexExpandingLookup && !element.wasResolved) return undefined; } } @@ -365,7 +458,7 @@ export class RequestParameterEvaluators { filters: lookup.filters.map((e) => visitExpr(mapOutputs, e, null)) }); - resolved = new TableValuedExpandingLookup(amountOfLookups++, (request) => [ + resolved = new TableValuedExpandingLookup(amountOfLookups++, lookup.outputs.length, (request) => [ ...filterParameterRows(prepared.evaluate(parametersForRequest(request, mapInputs.instantiation))) ]); } @@ -394,8 +487,6 @@ export type PreparedExpandingLookup = TableValuedExpandingLookup | ParameterInde type CloneParameter = (original: PreparedParameterValue) => PreparedParameterValue; abstract class BasePreparedExpandingLookup implements ResultSetElement { - wasResolved = false; - constructor(readonly resultSetIndex: number) {} abstract clone(parameters: CloneParameter): BasePreparedExpandingLookup; @@ -404,19 +495,21 @@ abstract class BasePreparedExpandingLookup implements ResultSetElement { class TableValuedExpandingLookup extends BasePreparedExpandingLookup { constructor( resultSetIndex: number, + readonly columnCount: number, readonly read: (request: RequestParameters) => SqliteParameterValue[][] ) { super(resultSetIndex); } override clone(): TableValuedExpandingLookup { - const lookup = new TableValuedExpandingLookup(this.resultSetIndex, this.read); - lookup.wasResolved = this.wasResolved; - return lookup; + // Immutable instance + return this; } } class ParameterIndexExpandingLookup extends BasePreparedExpandingLookup { + wasResolved: boolean = false; + constructor( resultSetIndex: number, readonly lookup: ParameterIndexLookupCreator, diff --git a/packages/sync-rules/test/src/sync_plan/evaluator/evaluator.test.ts b/packages/sync-rules/test/src/sync_plan/evaluator/evaluator.test.ts index 1b70011ed..ec6b6c7ce 100644 --- a/packages/sync-rules/test/src/sync_plan/evaluator/evaluator.test.ts +++ b/packages/sync-rules/test/src/sync_plan/evaluator/evaluator.test.ts @@ -1,6 +1,7 @@ import * as sqlite from 'node:sqlite'; import { describe, expect, test } from 'vitest'; import { + CompatibilityContext, DEFAULT_HYDRATION_STATE, DEFAULT_TAG, deserializeSyncPlan, @@ -550,7 +551,7 @@ streams: expect(querier.staticBuckets.map((e) => e.bucket)).toStrictEqual(['stream|0["user"]']); }); - syncTest('intersection of request data', ({ sync }) => { + syncTest('intersection of scalar and expanded request data', ({ sync }) => { const desc = sync.prepareSyncStreams(` config: edition: 3 @@ -576,6 +577,121 @@ streams: expect(queryWith('p1', ['p1', 'p2'])).toStrictEqual(['stream|0["p1"]']); }); + syncTest( + 'intersection of large expanded request data', + ({ sync }) => { + // Regression test for https://github.com/powersync-ja/powersync-service/pull/782#pullrequestreview-5140431191: + // For intersections between request parameters, ensure we use an efficient filter. This test relies on a 1s + // timeout, which catches quadratic behavior. + const desc = sync.prepareSyncStreams(` +config: + edition: 3 + +streams: + stream: + auto_subscribe: true + query: SELECT * FROM issues WHERE a IN auth.parameter('x') AND a IN auth.parameter('y') AND a IN auth.parameter('z') +`); + + const x: number[] = []; + const y: number[] = []; + const z: number[] = [1, 2, 3, 4, 5]; + + for (let i = 0; i < 50_000; i++) { + x.push(i); + y.push(i); + } + + const { querier, errors } = desc.getBucketParameterQuerier({ + globalParameters: requestParameters({ sub: 'user', x, y, z }), + hasDefaultStreams: true, + streams: {} + }); + expect(errors).toStrictEqual([]); + expect(querier.hasDynamicBuckets).toStrictEqual(false); + expect(querier.staticBuckets.map((b) => b.bucket)).toStrictEqual(z.map((param) => `stream|0[${param}]`)); + }, + 1000 + ); + + syncTest('partially overlapping intersections', ({ sync }) => { + // This query has two intersections: (x, y) and (x, z). Currently, the compiler emits distinct parameters for the + // same lookup. But this plan otherwise corresponds to the query SELECT * FROM issues WHERE + // a IN auth.parameter('x') AND a IN auth.parameter('y') AND + // b IN auth.parameter('x') AND b IN auth.parameter('z') + + function parameter(key: string) { + return { + type: 'table_valued', + functionName: 'json_each', + functionInputs: [ + { + type: 'function', + function: '->>', + parameters: [ + { type: 'data', source: { request: 'auth' } }, + { type: 'lit_string', value: key } + ] + } + ], + outputs: [{ type: 'data', source: { column: 'value' } }], + filters: [] + }; + } + const plan = deserializeSyncPlan({ + dataSources: [], + buckets: [{ hash: 5373205, uniqueName: 'stream|0', sources: [] }], + parameterIndexes: [], + streams: [ + { + stream: { name: 'stream', priority: 3, isSubscribedByDefault: true }, + queriers: [ + { + requestFilters: [], + lookupStages: [[parameter('x'), parameter('y'), parameter('z')]], + bucket: 0, + sourceInstantiation: [ + { + type: 'intersection', + values: [ + { type: 'lookup', lookup: { stageId: 0, idInStage: 0 }, resultIndex: 0 }, + { type: 'lookup', lookup: { stageId: 0, idInStage: 1 }, resultIndex: 0 } + ] + }, + { + type: 'intersection', + values: [ + { type: 'lookup', lookup: { stageId: 0, idInStage: 0 }, resultIndex: 0 }, + { type: 'lookup', lookup: { stageId: 0, idInStage: 2 }, resultIndex: 0 } + ] + } + ] + } + ] + } + ], + version: 1 + }); + const desc = sync.hydrateConfig( + new PrecompiledSyncConfig(plan, new CompatibilityContext({ edition: 3 }), { + defaultSchema: 'ignored', + sourceText: 'ignored' + }) + ); + + const { querier, errors } = desc.getBucketParameterQuerier({ + globalParameters: requestParameters({ sub: 'user', x: [1, 2, 3, 4, 5, 6], y: [2, 4, 6], z: [1, 3, 5] }), + hasDefaultStreams: true, + streams: {} + }); + + expect(errors).toStrictEqual([]); + expect(querier.hasDynamicBuckets).toStrictEqual(false); + + // There are no rows where x both intersects with y and z, so this should result in no buckets. + expect(querier.staticBuckets.map((b) => b.bucket)).toStrictEqual([]); + }); + syncTest('parameter lookups', async ({ sync }) => { const desc = sync.prepareSyncStreams(` config: diff --git a/packages/sync-rules/test/src/sync_plan/evaluator/utils.ts b/packages/sync-rules/test/src/sync_plan/evaluator/utils.ts index 08a06cffe..850a432c8 100644 --- a/packages/sync-rules/test/src/sync_plan/evaluator/utils.ts +++ b/packages/sync-rules/test/src/sync_plan/evaluator/utils.ts @@ -21,6 +21,7 @@ interface SyncTest { params?: HydrateSyncConfigParams, options?: PrepareStreamsOptions ): HydratedSyncConfig; + hydrateConfig(config: SyncConfig, params?: HydrateSyncConfigParams): HydratedSyncConfig; } export const syncTest = test.extend<{ sync: SyncTest }>({ @@ -34,10 +35,11 @@ export const syncTest = test.extend<{ sync: SyncTest }>({ return config; }, + hydrateConfig(config, params?: HydrateSyncConfigParams) { + return config.hydrate(params ?? { hydrationState: DEFAULT_HYDRATION_STATE, sqlite: nodeSqlite(sqlite) }); + }, prepareSyncStreams(inputs, params?: HydrateSyncConfigParams, options?: PrepareStreamsOptions) { - return this.prepareWithoutHydration(inputs, options).hydrate( - params ?? { hydrationState: DEFAULT_HYDRATION_STATE, sqlite: nodeSqlite(sqlite) } - ); + return this.hydrateConfig(this.prepareWithoutHydration(inputs, options), params); } }); } From 1183c12e160680e608c35a580cdc224bec30f27e Mon Sep 17 00:00:00 2001 From: Simon Binder Date: Mon, 14 Sep 2026 11:46:14 +0200 Subject: [PATCH 14/15] Simplify, AI feedback --- .../evaluator/parameter_evaluator.ts | 58 +++++++++---------- .../src/sync_plan/evaluator/evaluator.test.ts | 2 +- 2 files changed, 27 insertions(+), 33 deletions(-) diff --git a/packages/sync-rules/src/sync_plan/evaluator/parameter_evaluator.ts b/packages/sync-rules/src/sync_plan/evaluator/parameter_evaluator.ts index af0b2ced5..e152478ef 100644 --- a/packages/sync-rules/src/sync_plan/evaluator/parameter_evaluator.ts +++ b/packages/sync-rules/src/sync_plan/evaluator/parameter_evaluator.ts @@ -119,18 +119,14 @@ export class RequestParameterEvaluators { */ partiallyInstantiate(input: PartialInstantiationInput): SqliteParameterValue[][] | undefined { try { - const mappedIntersectionConstraints = new Set(); + const mappedIntersectionConstraints: IntersectionConstraint[] = []; // Optimization: For request-parameter based lookups with a single column (the most common kind), use a set to // immediately filter out rows that don't intersect with another request parameter. ResultSet.multiply does the // same thing, but that requires going through the combination of all rows whereas the de-duplication here is // linear. interface RequestParameterIntersection { - constraint: IntersectionConstraint; - // If the intersection consists entirely of request parameters, and all mentioned request parameters are in - // exactly one intersection, we don't need to register the constraint on the result set as it's covered by the - // optimization. - replacesConstraint: boolean; + fixedValue: SqliteParameterValue | undefined; amountOfRequestParameters: number; // Count in how many parameters a value was seen, we filter out values not present in all parameters. rows: Map; @@ -141,11 +137,9 @@ export class RequestParameterEvaluators { for (const constraint of this.intersectionConstraints) { const mapped = this.#createIntersectionResult(constraint, input); - mappedIntersectionConstraints.add(mapped); const parameterIntersection: RequestParameterIntersection = { - constraint: mapped, - replacesConstraint: true, + fixedValue: mapped.fixedValue, amountOfRequestParameters: 0, rows: new Map(), materializedRows: [] @@ -163,11 +157,6 @@ export class RequestParameterEvaluators { if (existing != null) { existing.push(parameterIntersection); - - // This parameter is part of multiple intersections. - for (const intersection of existing) { - intersection.replacesConstraint = false; - } } else { parameterToIntersection.set(column.lookup, [parameterIntersection]); } @@ -175,16 +164,9 @@ export class RequestParameterEvaluators { } if (parameterIntersection.amountOfRequestParameters < constraint.inputs.length) { - // This intersection doesn't entirely consist of request parameters. - parameterIntersection.replacesConstraint = false; - } - } - - for (const intersections of parameterToIntersection.values()) { - for (const intersection of intersections) { - if (intersection.replacesConstraint) { - mappedIntersectionConstraints.delete(intersection.constraint); - } + // If the intersection consists entirely of parameter lookups, it's fully covered by the optimization and + // doesn't need to be tracked in the result set. + mappedIntersectionConstraints.push(mapped); } } @@ -201,15 +183,15 @@ export class RequestParameterEvaluators { this.#checkInstantiable(); } else { // For this parameter to be part of the intersection optimization, it must return exactly one column. - for (const [value] of outputs) { + for (const [value] of new Set(outputs)) { for (const intersection of intersections) { - const fixed = intersection.constraint.fixedValue; + const fixed = intersection.fixedValue; if (fixed != null && value != fixed) continue; const matchedSources = (intersection.rows.get(value) ?? 0) + 1; intersection.rows.set(value, matchedSources); - if (matchedSources == intersection.amountOfRequestParameters) { + if (matchedSources === intersection.amountOfRequestParameters) { intersection.materializedRows.push([value]); } } @@ -226,12 +208,24 @@ export class RequestParameterEvaluators { } for (const [parameter, intersections] of parameterToIntersection.entries()) { - // Find the smallest intersection constraining this parameter, then instantiate the parameter to that. - const smallestIntersection = intersections.reduce((acc, intersection) => - intersection.materializedRows.length < acc.materializedRows.length ? intersection : acc - ); + if (intersections.length == 1) { + this.resultSet.multiply(parameter.resultSetIndex, intersections[0].materializedRows); + } else { + const rows = intersections.reduce((acc, intersection) => + intersection.materializedRows.length < acc.materializedRows.length ? intersection : acc + ).materializedRows; + + this.resultSet.multiply( + parameter.resultSetIndex, + rows.filter(([value]) => { + return intersections.every((intersection) => { + const matchingCount = intersection.rows.get(value); + return matchingCount === intersection.amountOfRequestParameters; + }); + }) + ); + } - this.resultSet.multiply(parameter.resultSetIndex, smallestIntersection.materializedRows); this.#checkInstantiable(); } diff --git a/packages/sync-rules/test/src/sync_plan/evaluator/evaluator.test.ts b/packages/sync-rules/test/src/sync_plan/evaluator/evaluator.test.ts index ec6b6c7ce..0dce5625d 100644 --- a/packages/sync-rules/test/src/sync_plan/evaluator/evaluator.test.ts +++ b/packages/sync-rules/test/src/sync_plan/evaluator/evaluator.test.ts @@ -573,7 +573,7 @@ streams: return querier.staticBuckets.map((e) => e.bucket); } - expect(queryWith('p1', ['p2'])).toStrictEqual([]); + expect(queryWith('p1', ['p2', 'p2'])).toStrictEqual([]); expect(queryWith('p1', ['p1', 'p2'])).toStrictEqual(['stream|0["p1"]']); }); From 6f398e7b3ec5646a656fc8aa9dfe2f5f34da89df Mon Sep 17 00:00:00 2001 From: Simon Binder Date: Mon, 14 Sep 2026 12:08:01 +0200 Subject: [PATCH 15/15] Use strict equality --- .../sync-rules/src/sync_plan/evaluator/parameter_evaluator.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/packages/sync-rules/src/sync_plan/evaluator/parameter_evaluator.ts b/packages/sync-rules/src/sync_plan/evaluator/parameter_evaluator.ts index e152478ef..cab853862 100644 --- a/packages/sync-rules/src/sync_plan/evaluator/parameter_evaluator.ts +++ b/packages/sync-rules/src/sync_plan/evaluator/parameter_evaluator.ts @@ -186,7 +186,7 @@ export class RequestParameterEvaluators { for (const [value] of new Set(outputs)) { for (const intersection of intersections) { const fixed = intersection.fixedValue; - if (fixed != null && value != fixed) continue; + if (fixed != null && value !== fixed) continue; const matchedSources = (intersection.rows.get(value) ?? 0) + 1;