diff --git a/.changeset/smart-jokes-walk.md b/.changeset/smart-jokes-walk.md new file mode 100644 index 000000000..07bea5169 --- /dev/null +++ b/.changeset/smart-jokes-walk.md @@ -0,0 +1,7 @@ +--- +'@powersync/service-core-tests': minor +'@powersync/service-core': minor +'@powersync/service-sync-rules': minor +--- + +[Internal] Track source on buckets and parameter indexes. diff --git a/packages/service-core-tests/src/tests/register-data-storage-parameter-tests.ts b/packages/service-core-tests/src/tests/register-data-storage-parameter-tests.ts index 2324f8f43..3a3d7d1d6 100644 --- a/packages/service-core-tests/src/tests/register-data-storage-parameter-tests.ts +++ b/packages/service-core-tests/src/tests/register-data-storage-parameter-tests.ts @@ -1,9 +1,8 @@ import { CURRENT_STORAGE_VERSION, JwtPayload, storage, updateSyncRulesFromYaml } from '@powersync/service-core'; import { RequestParameters, ScopedParameterLookup, SqliteJsonRow } from '@powersync/service-sync-rules'; -import { ParameterLookupScope } from '@powersync/service-sync-rules/src/HydrationState.js'; import { expect, test } from 'vitest'; import * as test_utils from '../test-utils/test-utils-index.js'; -import { bucketRequest } from './util.js'; +import { bucketRequest, parameterLookupScope } from './util.js'; /** * @example @@ -19,7 +18,7 @@ export function registerDataStorageParameterTests(config: storage.TestStorageCon const generateStorageFactory = config.factory; const storageVersion = config.storageVersion ?? CURRENT_STORAGE_VERSION; const TEST_TABLE = test_utils.makeTestTable('test', ['id'], config); - const MYBUCKET_1: ParameterLookupScope = { lookupName: 'mybucket', queryId: '1' }; + const MYBUCKET_1 = parameterLookupScope('mybucket', '1'); test('save and load parameters', async () => { await using factory = await generateStorageFactory(); @@ -375,7 +374,7 @@ bucket_definitions: const buckets = await querier.queryDynamicBucketDescriptions({ async getParameterSets(lookups) { - expect(lookups).toEqual([ScopedParameterLookup.direct({ lookupName: 'by_workspace', queryId: '1' }, ['u1'])]); + expect(lookups).toEqual([ScopedParameterLookup.direct(parameterLookupScope('by_workspace', '1'), ['u1'])]); const parameter_sets = await checkpoint.getParameterSets(lookups); expect(parameter_sets).toEqual([{ workspace_id: 'workspace1' }]); @@ -457,9 +456,7 @@ bucket_definitions: const buckets = await querier.queryDynamicBucketDescriptions({ async getParameterSets(lookups) { - expect(lookups).toEqual([ - ScopedParameterLookup.direct({ lookupName: 'by_public_workspace', queryId: '1' }, []) - ]); + expect(lookups).toEqual([ScopedParameterLookup.direct(parameterLookupScope('by_public_workspace', '1'), [])]); const parameter_sets = await checkpoint.getParameterSets(lookups); parameter_sets.sort((a, b) => JSON.stringify(a).localeCompare(JSON.stringify(b))); @@ -576,8 +573,8 @@ bucket_definitions: }) ).map((e) => e.bucket); expect(foundLookups).toEqual([ - ScopedParameterLookup.direct({ lookupName: 'by_workspace', queryId: '1' }, []), - ScopedParameterLookup.direct({ lookupName: 'by_workspace', queryId: '2' }, ['u1']) + ScopedParameterLookup.direct(parameterLookupScope('by_workspace', '1'), []), + ScopedParameterLookup.direct(parameterLookupScope('by_workspace', '2'), ['u1']) ]); parameter_sets.sort((a, b) => JSON.stringify(a).localeCompare(JSON.stringify(b))); expect(parameter_sets).toEqual([{ workspace_id: 'workspace1' }, { workspace_id: 'workspace3' }]); @@ -703,13 +700,7 @@ streams: const checkpoint = await bucketStorage.getCheckpoint(); const parameters = await checkpoint.getParameterSets([ - ScopedParameterLookup.direct( - { - lookupName: 'lookup', - queryId: '0' - }, - ['baz'] - ) + ScopedParameterLookup.direct(parameterLookupScope('lookup', '0'), ['baz']) ]); expect(parameters).toEqual([ { diff --git a/packages/service-core-tests/src/tests/register-parameter-compacting-tests.ts b/packages/service-core-tests/src/tests/register-parameter-compacting-tests.ts index 49f0a1381..0e6106ede 100644 --- a/packages/service-core-tests/src/tests/register-parameter-compacting-tests.ts +++ b/packages/service-core-tests/src/tests/register-parameter-compacting-tests.ts @@ -2,6 +2,7 @@ import { storage, updateSyncRulesFromYaml } from '@powersync/service-core'; import { ScopedParameterLookup } from '@powersync/service-sync-rules'; import { expect, test } from 'vitest'; import * as test_utils from '../test-utils/test-utils-index.js'; +import { parameterLookupScope } from './util.js'; export function registerParameterCompactTests(config: storage.TestStorageConfig) { const generateStorageFactory = config.factory; @@ -43,7 +44,7 @@ bucket_definitions: await batch.commit('1/1'); }); - const lookup = ScopedParameterLookup.direct({ lookupName: 'test', queryId: '1' }, ['t1']); + const lookup = ScopedParameterLookup.direct(parameterLookupScope('test', '1'), ['t1']); const checkpoint1 = await bucketStorage.getCheckpoint(); const parameters1 = await checkpoint1.getParameterSets([lookup]); @@ -155,7 +156,7 @@ bucket_definitions: await batch.commit('3/1'); }); - const lookup = ScopedParameterLookup.direct({ lookupName: 'test', queryId: '1' }, ['u1']); + const lookup = ScopedParameterLookup.direct(parameterLookupScope('test', '1'), ['u1']); const checkpoint1 = await bucketStorage.getCheckpoint(); const parameters1 = await checkpoint1.getParameterSets([lookup]); diff --git a/packages/service-core-tests/src/tests/util.ts b/packages/service-core-tests/src/tests/util.ts index 85e0e95c5..fcb0147ba 100644 --- a/packages/service-core-tests/src/tests/util.ts +++ b/packages/service-core-tests/src/tests/util.ts @@ -1,4 +1,11 @@ import { storage } from '@powersync/service-core'; +import { + ParameterIndexLookupCreator, + SourceTableInterface, + SqliteRow, + TablePattern +} from '@powersync/service-sync-rules'; +import { ParameterLookupScope } from '@powersync/service-sync-rules/src/HydrationState.js'; export function bucketRequest(syncRules: storage.PersistedSyncRulesContent, bucketName: string): string { if (/^\d+#/.test(bucketName)) { @@ -19,3 +26,30 @@ export function bucketRequestMap( ): Map { return new Map(Array.from(buckets, ([bucketName, opId]) => [bucketRequest(syncRules, bucketName), opId])); } + +const EMPTY_LOOKUP_SOURCE: ParameterIndexLookupCreator = { + get defaultLookupScope(): ParameterLookupScope { + return { + lookupName: 'lookup', + queryId: '0', + source: EMPTY_LOOKUP_SOURCE + }; + }, + getSourceTables(): Set { + return new Set(); + }, + evaluateParameterRow(_sourceTable: SourceTableInterface, _row: SqliteRow) { + return []; + }, + tableSyncsParameters(_table: SourceTableInterface): boolean { + return false; + } +}; + +export function parameterLookupScope( + lookupName: string, + queryId: string, + source: ParameterIndexLookupCreator = EMPTY_LOOKUP_SOURCE +): ParameterLookupScope { + return { lookupName, queryId, source }; +} diff --git a/packages/service-core/src/sync/BucketChecksumState.ts b/packages/service-core/src/sync/BucketChecksumState.ts index 6c3f780c8..b652e9f84 100644 --- a/packages/service-core/src/sync/BucketChecksumState.ts +++ b/packages/service-core/src/sync/BucketChecksumState.ts @@ -7,7 +7,8 @@ import { QuerierError, RequestedStream, RequestParameters, - ResolvedBucket + ResolvedBucket, + mergeBuckets } from '@powersync/service-sync-rules'; import * as storage from '../storage/storage-index.js'; @@ -168,7 +169,7 @@ export class BucketChecksumState { } // Subset of buckets for which there may be new data in this batch. - let bucketsToFetch: BucketDescription[]; + let bucketsToFetch: ResolvedBucket[]; let checkpointLine: util.StreamingSyncCheckpointDiff | util.StreamingSyncCheckpoint; @@ -207,10 +208,7 @@ export class BucketChecksumState { ...this.parameterState.translateResolvedBucket(bucketDescriptionMap.get(e.bucket)!, streamNameToIndex) })); bucketsToFetch = [...generateBucketsToFetch].map((b) => { - return { - priority: bucketDescriptionMap.get(b)!.priority, - bucket: b - }; + return bucketDescriptionMap.get(b)!; }); deferredLog = () => { @@ -265,7 +263,7 @@ export class BucketChecksumState { totalParamResults ); }; - bucketsToFetch = allBuckets.map((b) => ({ bucket: b.bucket, priority: b.priority })); + bucketsToFetch = allBuckets; const subscriptions: util.StreamDescription[] = []; const streamNameToIndex = new Map(); @@ -342,7 +340,7 @@ export class BucketChecksumState { deferredLog(); }, - getFilteredBucketPositions: (buckets?: BucketDescription[]): Map => { + getFilteredBucketPositions: (buckets?: ResolvedBucket[]): Map => { if (!hasAdvanced) { throw new ServiceAssertionError('Call line.advance() before getFilteredBucketPositions()'); } @@ -660,7 +658,7 @@ export class BucketParameterState { export interface CheckpointLine { checkpointLine: util.StreamingSyncCheckpointDiff | util.StreamingSyncCheckpoint; - bucketsToFetch: BucketDescription[]; + bucketsToFetch: ResolvedBucket[]; /** * Call when a checkpoint line is being sent to a client, to update the internal state. @@ -672,7 +670,7 @@ export interface CheckpointLine { * * @param bucketsToFetch List of buckets to fetch - either this.bucketsToFetch, or a subset of it. Defaults to this.bucketsToFetch. */ - getFilteredBucketPositions(bucketsToFetch?: BucketDescription[]): Map; + getFilteredBucketPositions(bucketsToFetch?: ResolvedBucket[]): Map; /** * Update the position of bucket data the client has, after it was sent to the client. @@ -762,32 +760,3 @@ function limitedBuckets(buckets: string[] | { bucket: string }[], limit: number) const limited = buckets.slice(0, limit); return `${JSON.stringify(limited)}...`; } - -/** - * Resolves duplicate buckets in the given array, merging the inclusion reasons for duplicate. - * - * It's possible for duplicates to occur when a stream has multiple subscriptions, consider e.g. - * - * ``` - * sync_streams: - * assets_by_category: - * query: select * from assets where category in (request.parameters() -> 'categories') - * ``` - * - * Here, a client might subscribe once with `{"categories": [1]}` and once with `{"categories": [1, 2]}`. Since each - * subscription is evaluated independently, this would lead to three buckets, with a duplicate `assets_by_category[1]` - * bucket. - */ -function mergeBuckets(buckets: ResolvedBucket[]): ResolvedBucket[] { - const byBucketId: Record = {}; - - for (const bucket of buckets) { - if (Object.hasOwn(byBucketId, bucket.bucket)) { - byBucketId[bucket.bucket].inclusion_reasons.push(...bucket.inclusion_reasons); - } else { - byBucketId[bucket.bucket] = structuredClone(bucket); - } - } - - return Object.values(byBucketId); -} diff --git a/packages/service-core/src/sync/sync.ts b/packages/service-core/src/sync/sync.ts index 24ebdfd06..b2ea9c61e 100644 --- a/packages/service-core/src/sync/sync.ts +++ b/packages/service-core/src/sync/sync.ts @@ -1,5 +1,11 @@ import { JSONBig, JsonContainer } from '@powersync/service-jsonbig'; -import { BucketDescription, BucketPriority, HydratedSyncRules, SqliteJsonValue } from '@powersync/service-sync-rules'; +import { + BucketDescription, + BucketPriority, + HydratedSyncRules, + ResolvedBucket, + SqliteJsonValue +} from '@powersync/service-sync-rules'; import { AbortError } from 'ix/aborterror.js'; @@ -179,7 +185,7 @@ async function* streamResponseInner( // receive a sync complete message after the synchronization is done (which happens in the last // bucketDataInBatches iteration). Without any batch, the line is missing and clients might not complete their // sync properly. - const priorityBatches: [BucketPriority | null, BucketDescription[]][] = bucketsByPriority; + const priorityBatches: [BucketPriority | null, ResolvedBucket[]][] = bucketsByPriority; if (priorityBatches.length == 0) { priorityBatches.push([null, []]); } @@ -257,7 +263,7 @@ interface BucketDataRequest { /** Contains current bucket state. Modified by the request as data is sent. */ checkpointLine: CheckpointLine; /** Subset of checkpointLine.bucketsToFetch, filtered by priority. */ - bucketsToFetch: BucketDescription[]; + bucketsToFetch: ResolvedBucket[]; /** Whether data lines should be encoded in a legacy format where {@link util.OplogEntry.data} is a nested object. */ legacyDataLines: boolean; /** Signals that the connection was aborted and that streaming should stop ASAP. */ diff --git a/packages/service-core/test/src/sync/BucketChecksumState.test.ts b/packages/service-core/test/src/sync/BucketChecksumState.test.ts index f706bc018..a59f6a223 100644 --- a/packages/service-core/test/src/sync/BucketChecksumState.test.ts +++ b/packages/service-core/test/src/sync/BucketChecksumState.test.ts @@ -14,15 +14,43 @@ import { } from '@/index.js'; import { JSONBig } from '@powersync/service-jsonbig'; import { + ParameterIndexLookupCreator, RequestJwtPayload, ScopedParameterLookup, SqliteJsonRow, + SqliteRow, SqlSyncRules, + TablePattern, + SourceTableInterface, versionedHydrationState } from '@powersync/service-sync-rules'; +import { ParameterLookupScope } from '@powersync/service-sync-rules/src/HydrationState.js'; import { beforeEach, describe, expect, test } from 'vitest'; describe('BucketChecksumState', () => { + const LOOKUP_SOURCE: ParameterIndexLookupCreator = { + get defaultLookupScope(): ParameterLookupScope { + return { + lookupName: 'lookup', + queryId: '0', + source: LOOKUP_SOURCE + }; + }, + getSourceTables(): Set { + return new Set(); + }, + evaluateParameterRow(_sourceTable: SourceTableInterface, _row: SqliteRow) { + return []; + }, + tableSyncsParameters(_table: SourceTableInterface): boolean { + return false; + } + }; + + function lookupScope(lookupName: string, queryId: string): ParameterLookupScope { + return { lookupName, queryId, source: LOOKUP_SOURCE }; + } + // Single global[] bucket. // We don't care about data in these tests const SYNC_RULES_GLOBAL = SqlSyncRules.fromYaml( @@ -94,7 +122,7 @@ bucket_definitions: streams: [{ name: 'global', is_default: true, errors: [] }] } }); - expect(line.bucketsToFetch).toEqual([ + expect(line.bucketsToFetch).toMatchObject([ { bucket: '1#global[]', priority: 3 @@ -166,7 +194,7 @@ bucket_definitions: streams: [{ name: 'global', is_default: true, errors: [] }] } }); - expect(line.bucketsToFetch).toEqual([ + expect(line.bucketsToFetch).toMatchObject([ { bucket: '1#global[]', priority: 3 @@ -206,7 +234,7 @@ bucket_definitions: streams: [{ name: 'global', is_default: true, errors: [] }] } }); - expect(line.bucketsToFetch).toEqual([ + expect(line.bucketsToFetch).toMatchObject([ { bucket: '2#global[1]', priority: 3 @@ -274,7 +302,7 @@ bucket_definitions: streams: [{ name: 'global', is_default: true, errors: [] }] } }); - expect(line.bucketsToFetch).toEqual([ + expect(line.bucketsToFetch).toMatchObject([ { bucket: '1#global[]', priority: 3 @@ -337,7 +365,7 @@ bucket_definitions: write_checkpoint: undefined } }); - expect(line2.bucketsToFetch).toEqual([{ bucket: '2#global[1]', priority: 3 }]); + expect(line2.bucketsToFetch).toMatchObject([{ bucket: '2#global[1]', priority: 3 }]); }); test('invalidating all buckets', async () => { @@ -387,7 +415,7 @@ bucket_definitions: write_checkpoint: undefined } }); - expect(line2.bucketsToFetch).toEqual([ + expect(line2.bucketsToFetch).toMatchObject([ { bucket: '2#global[1]', priority: 3 }, { bucket: '2#global[2]', priority: 3 } ]); @@ -424,7 +452,7 @@ bucket_definitions: streams: [{ name: 'global', is_default: true, errors: [] }] } }); - expect(line.bucketsToFetch).toEqual([ + expect(line.bucketsToFetch).toMatchObject([ { bucket: '2#global[1]', priority: 3 @@ -477,7 +505,7 @@ bucket_definitions: } }); // This should contain both buckets, even though only one changed. - expect(line2.bucketsToFetch).toEqual([ + expect(line2.bucketsToFetch).toMatchObject([ { bucket: '2#global[1]', priority: 3 @@ -513,7 +541,7 @@ bucket_definitions: const line = (await state.buildNextCheckpointLine({ base: storage.makeCheckpoint(1n, (lookups) => { - expect(lookups).toEqual([ScopedParameterLookup.direct({ lookupName: 'by_project', queryId: '1' }, ['u1'])]); + expect(lookups).toEqual([ScopedParameterLookup.direct(lookupScope('by_project', '1'), ['u1'])]); return [{ id: 1 }, { id: 2 }]; }), writeCheckpoint: null, @@ -548,7 +576,7 @@ bucket_definitions: write_checkpoint: undefined } }); - expect(line.bucketsToFetch).toEqual([ + expect(line.bucketsToFetch).toMatchObject([ { bucket: '3#by_project[1]', priority: 3 @@ -574,7 +602,7 @@ bucket_definitions: // Now we get a new line const line2 = (await state.buildNextCheckpointLine({ base: storage.makeCheckpoint(2n, (lookups) => { - expect(lookups).toEqual([ScopedParameterLookup.direct({ lookupName: 'by_project', queryId: '1' }, ['u1'])]); + expect(lookups).toEqual([ScopedParameterLookup.direct(lookupScope('by_project', '1'), ['u1'])]); return [{ id: 1 }, { id: 2 }, { id: 3 }]; }), writeCheckpoint: null, diff --git a/packages/sync-rules/src/BaseSqlDataQuery.ts b/packages/sync-rules/src/BaseSqlDataQuery.ts index 9d86a157d..caa962282 100644 --- a/packages/sync-rules/src/BaseSqlDataQuery.ts +++ b/packages/sync-rules/src/BaseSqlDataQuery.ts @@ -1,9 +1,9 @@ import { SelectedColumn } from 'pgsql-ast-parser'; +import { idFromData } from './cast.js'; import { SqlRuleError } from './errors.js'; import { ColumnDefinition } from './ExpressionType.js'; import { SourceTableInterface } from './SourceTableInterface.js'; import { AvailableTable, SqlTools } from './sql_filters.js'; -import { castAsText } from './sql_functions.js'; import { TablePattern } from './TablePattern.js'; import { QueryParameters, @@ -15,7 +15,7 @@ import { SqliteJsonRow, SqliteRow } from './types.js'; -import { filterJsonRow, idFromData } from './utils.js'; +import { filterJsonRow } from './utils.js'; export interface RowValueExtractor { extract(tables: QueryParameters, into: SqliteRow): void; diff --git a/packages/sync-rules/src/BucketDescription.ts b/packages/sync-rules/src/BucketDescription.ts index 8dd732f34..0fd5833ca 100644 --- a/packages/sync-rules/src/BucketDescription.ts +++ b/packages/sync-rules/src/BucketDescription.ts @@ -1,3 +1,5 @@ +import { BucketDataSource } from './BucketSource.js'; + /** * The priority in which to synchronize buckets. * @@ -19,7 +21,27 @@ export const isValidPriority = (i: number): i is BucketPriority => { return Number.isInteger(i) && i >= 0 && i <= 3; }; -export interface BucketDescription { +/** + * There is no _direct_ way to define that a property is not enumerable in TypeScript. + * + * A getter on a class does that indirectly. + * + * This is _not_ the same as defining `get source(): BucketDataSource` directly on the interface. + * + * We never instantiate or extend this class directly - we only use the type. + */ +abstract class NonEnumerableSourceClass { + private constructor() {} + + /** + * This is specifically not enumerable - must be excluded from tests and serialization. + */ + abstract get source(): BucketDataSource; +} + +export type NonEnumerableBucketDataSource = NonEnumerableSourceClass; + +export interface BucketDescription extends NonEnumerableBucketDataSource { /** * The id of the bucket, which is derived from the name of the bucket's definition * in the sync rules as well as the values returned by the parameter queries. diff --git a/packages/sync-rules/src/BucketParameterQuerier.ts b/packages/sync-rules/src/BucketParameterQuerier.ts index 4d603b0a8..de53c482c 100644 --- a/packages/sync-rules/src/BucketParameterQuerier.ts +++ b/packages/sync-rules/src/BucketParameterQuerier.ts @@ -1,5 +1,6 @@ import { JSONBig } from '@powersync/service-jsonbig'; import { ResolvedBucket } from './BucketDescription.js'; +import { ParameterIndexLookupCreator } from './BucketSource.js'; import { ParameterLookupScope } from './HydrationState.js'; import { RequestedStream } from './SqlSyncRules.js'; import { RequestParameters, SqliteJsonRow, SqliteJsonValue } from './types.js'; @@ -106,6 +107,7 @@ export function mergeBucketParameterQueriers(queriers: BucketParameterQuerier[]) export class ScopedParameterLookup { // bucket definition name, parameter query index, ...lookup values readonly values: readonly SqliteJsonValue[]; + readonly #source: ParameterIndexLookupCreator; #cachedSerializedForm?: string; @@ -119,22 +121,34 @@ export class ScopedParameterLookup { } static normalized(scope: ParameterLookupScope, lookup: UnscopedParameterLookup): ScopedParameterLookup { - return new ScopedParameterLookup([scope.lookupName, scope.queryId, ...lookup.lookupValues]); + return new ScopedParameterLookup(scope.source, [scope.lookupName, scope.queryId, ...lookup.lookupValues]); } /** * Primarily for test fixtures. */ static direct(scope: ParameterLookupScope, values: SqliteJsonValue[]): ScopedParameterLookup { - return new ScopedParameterLookup([scope.lookupName, scope.queryId, ...values.map(normalizeParameterValue)]); + return new ScopedParameterLookup(scope.source, [ + scope.lookupName, + scope.queryId, + ...values.map(normalizeParameterValue) + ]); } /** * * @param values must be pre-normalized (any integer converted into bigint) */ - private constructor(values: SqliteJsonValue[]) { + private constructor(source: ParameterIndexLookupCreator, values: SqliteJsonValue[]) { this.values = Object.freeze(values); + this.#source = source; + } + + /** + * Source, not enumerable (not checked in tests). + */ + get source() { + return this.#source; } } diff --git a/packages/sync-rules/src/BucketSource.ts b/packages/sync-rules/src/BucketSource.ts index 482fb5f5b..02e8c02bf 100644 --- a/packages/sync-rules/src/BucketSource.ts +++ b/packages/sync-rules/src/BucketSource.ts @@ -1,8 +1,8 @@ import { BucketParameterQuerier, - UnscopedParameterLookup, PendingQueriers, - ScopedParameterLookup + ScopedParameterLookup, + UnscopedParameterLookup } from './BucketParameterQuerier.js'; import { ColumnDefinition } from './ExpressionType.js'; import { DEFAULT_HYDRATION_STATE, HydrationState, ParameterLookupScope } from './HydrationState.js'; @@ -10,18 +10,18 @@ import { SourceTableInterface } from './SourceTableInterface.js'; import { GetQuerierOptions } from './SqlSyncRules.js'; import { TablePattern } from './TablePattern.js'; import { + EvaluatedParameters, EvaluatedParametersResult, EvaluatedRow, EvaluateRowOptions, EvaluationResult, isEvaluationError, - UnscopedEvaluationResult, SourceSchema, SqliteRow, UnscopedEvaluatedParametersResult, - EvaluatedParameters + UnscopedEvaluationResult } from './types.js'; -import { buildBucketName } from './utils.js'; +import { withBucketSource } from './utils.js'; export interface CreateSourceParams { hydrationState: HydrationState; @@ -171,12 +171,16 @@ export function hydrateEvaluateRow(hydrationState: HydrationState, source: Bucke if (isEvaluationError(result)) { return result; } - return { - bucket: buildBucketName(scope, result.serializedBucketParameters), - id: result.id, - table: result.table, - data: result.data - } satisfies EvaluatedRow; + const evaluated: EvaluatedRow = withBucketSource( + { + bucket: scope.bucketPrefix + result.serializedBucketParameters, + id: result.id, + table: result.table, + data: result.data + }, + scope.source + ); + return evaluated; }); }; } diff --git a/packages/sync-rules/src/HydrationState.ts b/packages/sync-rules/src/HydrationState.ts index f836a62b4..996de8056 100644 --- a/packages/sync-rules/src/HydrationState.ts +++ b/packages/sync-rules/src/HydrationState.ts @@ -3,12 +3,16 @@ import { BucketDataSource, ParameterIndexLookupCreator } from './BucketSource.js export interface BucketDataScope { /** The prefix is the bucket name before the parameters. */ bucketPrefix: string; + /** Source used to generate buckets. */ + source: BucketDataSource; } export interface ParameterLookupScope { /** The lookup name + queryid is used to reference the parameter lookup record. */ lookupName: string; queryId: string; + /** Source used to generate parameter lookups. */ + source: ParameterIndexLookupCreator; } /** @@ -37,7 +41,8 @@ export interface HydrationState { export const DEFAULT_HYDRATION_STATE: HydrationState = { getBucketSourceScope(source: BucketDataSource) { return { - bucketPrefix: source.uniqueName + bucketPrefix: source.uniqueName, + source }; }, getParameterIndexLookupScope(source) { @@ -61,7 +66,8 @@ export function versionedHydrationState(version: number): HydrationState { return { getBucketSourceScope(source: BucketDataSource): BucketDataScope { return { - bucketPrefix: `${version}#${source.uniqueName}` + bucketPrefix: `${version}#${source.uniqueName}`, + source }; }, diff --git a/packages/sync-rules/src/SqlParameterQuery.ts b/packages/sync-rules/src/SqlParameterQuery.ts index 2f4cb8e8b..ea214e0bd 100644 --- a/packages/sync-rules/src/SqlParameterQuery.ts +++ b/packages/sync-rules/src/SqlParameterQuery.ts @@ -18,6 +18,7 @@ import { BucketDataSource, BucketParameterQuerierSource, GetQuerierOptions, + resolvedBucket, ScopedParameterLookup, UnscopedEvaluatedParameters, UnscopedEvaluatedParametersResult @@ -42,7 +43,7 @@ import { SqliteRow } from './types.js'; import { - buildBucketName, + bucketDescription, filterJsonRow, isJsonValue, isSelectStatement, @@ -337,7 +338,8 @@ export class SqlParameterQuery implements ParameterIndexLookupCreator { public get defaultLookupScope(): ParameterLookupScope { return { lookupName: this.descriptorName, - queryId: this.queryId + queryId: this.queryId, + source: this }; } @@ -440,11 +442,7 @@ export class SqlParameterQuery implements ParameterIndexLookupCreator { } const serializedParameters = serializeBucketParameters(this.bucketParameters, result); - - return { - bucket: buildBucketName(bucketScope, serializedParameters), - priority: this.priority - }; + return bucketDescription(bucketScope, serializedParameters, this.priority); }) .filter((lookup) => lookup != null); } @@ -547,11 +545,9 @@ export class SqlParameterQuery implements ParameterIndexLookupCreator { hasDynamicBuckets: true, queryDynamicBucketDescriptions: async (source: ParameterLookupSource) => { const bucketParameters = await source.getParameterSets(lookups); - return this.resolveBucketDescriptions(bucketParameters, requestParameters, bucketDataScope).map((bucket) => ({ - ...bucket, - definition: this.descriptorName, - inclusion_reasons: reasons - })); + return this.resolveBucketDescriptions(bucketParameters, requestParameters, bucketDataScope).map((bucket) => { + return resolvedBucket(bucket, { definition: this.descriptorName, inclusion_reasons: reasons }); + }); } }; } diff --git a/packages/sync-rules/src/StaticSqlParameterQuery.ts b/packages/sync-rules/src/StaticSqlParameterQuery.ts index f08422ee0..a5cdc6b7b 100644 --- a/packages/sync-rules/src/StaticSqlParameterQuery.ts +++ b/packages/sync-rules/src/StaticSqlParameterQuery.ts @@ -1,16 +1,16 @@ import { SelectFromStatement } from 'pgsql-ast-parser'; -import { BucketDescription, BucketPriority, DEFAULT_BUCKET_PRIORITY, ResolvedBucket } from './BucketDescription.js'; +import { BucketDescription, BucketPriority, DEFAULT_BUCKET_PRIORITY } from './BucketDescription.js'; import { BucketParameterQuerier, PendingQueriers } from './BucketParameterQuerier.js'; import { CreateSourceParams } from './BucketSource.js'; import { SqlRuleError } from './errors.js'; import { BucketDataScope } from './HydrationState.js'; -import { BucketDataSource, BucketParameterQuerierSource, GetQuerierOptions } from './index.js'; +import { BucketDataSource, BucketParameterQuerierSource, GetQuerierOptions, resolvedBucket } from './index.js'; import { SourceTableInterface } from './SourceTableInterface.js'; import { AvailableTable, SqlTools } from './sql_filters.js'; import { checkUnsupportedFeatures, isClauseError, sqliteBool } from './sql_support.js'; import { TablePattern } from './TablePattern.js'; import { ParameterValueClause, QueryParseOptions, RequestParameters, SqliteJsonValue } from './types.js'; -import { buildBucketName, isJsonValue, serializeBucketParameters } from './utils.js'; +import { bucketDescription, isJsonValue, serializeBucketParameters } from './utils.js'; import { DetectRequestParameters } from './validators.js'; export interface StaticSqlParameterQueryOptions { @@ -181,11 +181,7 @@ export class StaticSqlParameterQuery { return { pushBucketParameterQueriers: (result: PendingQueriers, options: GetQuerierOptions) => { const staticBuckets = this.getStaticBucketDescriptions(options.globalParameters, bucketScope).map((desc) => { - return { - ...desc, - definition: this.descriptorName, - inclusion_reasons: ['default'] - } satisfies ResolvedBucket; + return resolvedBucket(desc, { definition: this.descriptorName, inclusion_reasons: ['default'] }); }); if (staticBuckets.length == 0) { @@ -224,13 +220,9 @@ export class StaticSqlParameterQuery { } const serializedParamters = serializeBucketParameters(this.bucketParameters, result); + const desc = bucketDescription(bucketSourceScope, serializedParamters, this.priority); - return [ - { - bucket: buildBucketName(bucketSourceScope, serializedParamters), - priority: this.priority - } - ]; + return [desc]; } get hasAuthenticatedBucketParameters(): boolean { diff --git a/packages/sync-rules/src/TableValuedFunctionSqlParameterQuery.ts b/packages/sync-rules/src/TableValuedFunctionSqlParameterQuery.ts index 8e176bfed..29a07d7cd 100644 --- a/packages/sync-rules/src/TableValuedFunctionSqlParameterQuery.ts +++ b/packages/sync-rules/src/TableValuedFunctionSqlParameterQuery.ts @@ -1,5 +1,5 @@ import { FromCall, SelectFromStatement } from 'pgsql-ast-parser'; -import { BucketDescription, BucketPriority, DEFAULT_BUCKET_PRIORITY, ResolvedBucket } from './BucketDescription.js'; +import { BucketDescription, BucketPriority, DEFAULT_BUCKET_PRIORITY } from './BucketDescription.js'; import { CreateSourceParams } from './BucketSource.js'; import { SqlRuleError } from './errors.js'; import { BucketDataScope } from './HydrationState.js'; @@ -23,7 +23,7 @@ import { SqliteJsonValue, SqliteRow } from './types.js'; -import { buildBucketName, isJsonValue, serializeBucketParameters } from './utils.js'; +import { bucketDescription, isJsonValue, resolvedBucket, serializeBucketParameters } from './utils.js'; import { DetectRequestParameters } from './validators.js'; export interface TableValuedFunctionSqlParameterQueryOptions { @@ -237,11 +237,7 @@ export class TableValuedFunctionSqlParameterQuery { return { pushBucketParameterQueriers: (result: PendingQueriers, options: GetQuerierOptions) => { const staticBuckets = this.getStaticBucketDescriptions(options.globalParameters, bucketScope).map((desc) => { - return { - ...desc, - definition: this.descriptorName, - inclusion_reasons: ['default'] - } satisfies ResolvedBucket; + return resolvedBucket(desc, { definition: this.descriptorName, inclusion_reasons: ['default'] }); }); if (staticBuckets.length == 0) { @@ -306,11 +302,7 @@ export class TableValuedFunctionSqlParameterQuery { } const serializedBucketParameters = serializeBucketParameters(this.bucketParameters, result); - - return { - bucket: buildBucketName(bucketScope, serializedBucketParameters), - priority: this.priority - }; + return bucketDescription(bucketScope, serializedBucketParameters, this.priority); } private visitParameterExtractorsAndCallClause(): DetectRequestParameters { diff --git a/packages/sync-rules/src/cast.ts b/packages/sync-rules/src/cast.ts new file mode 100644 index 000000000..5cc3dfcdf --- /dev/null +++ b/packages/sync-rules/src/cast.ts @@ -0,0 +1,123 @@ +import type { SqliteJsonRow, SqliteValue } from './types.js'; + +/** + * Extracts and normalizes the ID column from a row. + */ +export function idFromData(data: SqliteJsonRow): string { + let id = data.id; + if (typeof id != 'string') { + // While an explicit cast would be better, this covers against very common + // issues when initially testing out sync, for example when the id column is an + // auto-incrementing integer. + // If there is no id column, we use a blank id. This will result in the user syncing + // a single arbitrary row for this table - better than just not being able to sync + // anything. + id = castAsText(id) ?? ''; + } + return id; +} + +export const CAST_TYPES = new Set(['text', 'numeric', 'integer', 'real', 'blob']); + +const textEncoder = new TextEncoder(); +const textDecoder = new TextDecoder(); + +export function castAsText(value: SqliteValue): string | null { + if (value == null) { + return null; + } else if (value instanceof Uint8Array) { + return textDecoder.decode(value); + } else { + return value.toString(); + } +} + +export function castAsBlob(value: SqliteValue): Uint8Array | null { + if (value == null) { + return null; + } else if (value instanceof Uint8Array) { + return value; + } + + if (typeof value != 'string') { + value = value.toString(); + } + return textEncoder.encode(value); +} + +export function cast(value: SqliteValue, to: string) { + if (value == null) { + return null; + } + if (to == 'text') { + return castAsText(value); + } else if (to == 'numeric') { + if (value instanceof Uint8Array) { + value = textDecoder.decode(value); + } + if (typeof value == 'string') { + return parseNumeric(value); + } else if (typeof value == 'number' || typeof value == 'bigint') { + return value; + } else { + return 0n; + } + } else if (to == 'real') { + if (value instanceof Uint8Array) { + value = textDecoder.decode(value); + } + if (typeof value == 'string') { + const nr = parseFloat(value); + if (isNaN(nr)) { + return 0.0; + } else { + return nr; + } + } else if (typeof value == 'number') { + return value; + } else if (typeof value == 'bigint') { + return Number(value); + } else { + return 0.0; + } + } else if (to == 'integer') { + if (value instanceof Uint8Array) { + value = textDecoder.decode(value); + } + if (typeof value == 'string') { + return parseBigInt(value); + } else if (typeof value == 'number') { + return Number.isInteger(value) ? BigInt(value) : BigInt(Math.floor(value)); + } else if (typeof value == 'bigint') { + return value; + } else { + return 0n; + } + } else if (to == 'blob') { + return castAsBlob(value); + } else { + throw new Error(`Type not supported for cast: '${to}'`); + } +} + +function parseNumeric(text: string): bigint | number { + const match = /^\s*(\d+)(\.\d*)?(e[+\-]?\d+)?/i.exec(text); + if (!match) { + return 0n; + } + + if (match[2] != null || match[3] != null) { + const v = parseFloat(match[0]); + return isNaN(v) ? 0n : v; + } else { + return BigInt(match[1]); + } +} + +function parseBigInt(text: string): bigint { + const match = /^\s*(\d+)/.exec(text); + if (!match) { + return 0n; + } + return BigInt(match[1]); +} diff --git a/packages/sync-rules/src/compiler/sqlite.ts b/packages/sync-rules/src/compiler/sqlite.ts index 5ab0040f2..19ea4ac94 100644 --- a/packages/sync-rules/src/compiler/sqlite.ts +++ b/packages/sync-rules/src/compiler/sqlite.ts @@ -9,7 +9,7 @@ import { PGNode, SelectFromStatement } from 'pgsql-ast-parser'; -import { CAST_TYPES } from '../sql_functions.js'; +import { CAST_TYPES } from '../cast.js'; import { ColumnInRow, ConnectionParameter, diff --git a/packages/sync-rules/src/index.ts b/packages/sync-rules/src/index.ts index 3c68e85e2..420b459bc 100644 --- a/packages/sync-rules/src/index.ts +++ b/packages/sync-rules/src/index.ts @@ -28,6 +28,7 @@ export * from './types.js'; export * from './types/custom_sqlite_value.js'; export * from './types/time.js'; export * from './utils.js'; +export * from './cast.js'; export * from './HydrationState.js'; export * from './HydratedSyncRules.js'; diff --git a/packages/sync-rules/src/sql_functions.ts b/packages/sync-rules/src/sql_functions.ts index bbf89b662..f4f356691 100644 --- a/packages/sync-rules/src/sql_functions.ts +++ b/packages/sync-rules/src/sql_functions.ts @@ -10,6 +10,7 @@ import { ExpressionType, SqliteType, SqliteValueType, TYPE_INTEGER } from './Exp import * as uuid from 'uuid'; import { CustomSqliteValue } from './types/custom_sqlite_value.js'; import { CompatibilityContext, CompatibilityOption } from './compatibility.js'; +import { cast, CAST_TYPES, castAsBlob, castAsText } from './cast.js'; export const BASIC_OPERATORS = new Set([ '=', @@ -526,89 +527,6 @@ export function generateSqlFunctions(compatibility: CompatibilityContext) { }; } -export const CAST_TYPES = new Set(['text', 'numeric', 'integer', 'real', 'blob']); - -const textEncoder = new TextEncoder(); -const textDecoder = new TextDecoder(); - -export function castAsText(value: SqliteValue): string | null { - if (value == null) { - return null; - } else if (value instanceof Uint8Array) { - return textDecoder.decode(value); - } else { - return value.toString(); - } -} - -export function castAsBlob(value: SqliteValue): Uint8Array | null { - if (value == null) { - return null; - } else if (value instanceof Uint8Array) { - return value!; - } - - if (typeof value != 'string') { - value = value.toString(); - } - return textEncoder.encode(value); -} - -export function cast(value: SqliteValue, to: string) { - if (value == null) { - return null; - } - if (to == 'text') { - return castAsText(value); - } else if (to == 'numeric') { - if (value instanceof Uint8Array) { - value = textDecoder.decode(value); - } - if (typeof value == 'string') { - return parseNumeric(value); - } else if (typeof value == 'number' || typeof value == 'bigint') { - return value; - } else { - return 0n; - } - } else if (to == 'real') { - if (value instanceof Uint8Array) { - value = textDecoder.decode(value); - } - if (typeof value == 'string') { - const nr = parseFloat(value); - if (isNaN(nr)) { - return 0.0; - } else { - return nr; - } - } else if (typeof value == 'number') { - return value; - } else if (typeof value == 'bigint') { - return Number(value); - } else { - return 0.0; - } - } else if (to == 'integer') { - if (value instanceof Uint8Array) { - value = textDecoder.decode(value); - } - if (typeof value == 'string') { - return parseBigInt(value); - } else if (typeof value == 'number') { - return Number.isInteger(value) ? BigInt(value) : BigInt(Math.floor(value)); - } else if (typeof value == 'bigint') { - return value; - } else { - return 0n; - } - } else if (to == 'blob') { - return castAsBlob(value); - } else { - throw new Error(`Type not supported for cast: '${to}'`); - } -} - export function sqliteTypeOf(arg: SqliteInputValue): SqliteValueType { if (arg == null) { return 'null'; @@ -644,28 +562,6 @@ export function parseGeometry(value?: SqliteValue) { return geo; } -function parseNumeric(text: string): bigint | number { - const match = /^\s*(\d+)(\.\d*)?(e[+\-]?\d+)?/i.exec(text); - if (!match) { - return 0n; - } - - if (match[2] != null || match[3] != null) { - const v = parseFloat(match[0]); - return isNaN(v) ? 0n : v; - } else { - return BigInt(match[1]); - } -} - -function parseBigInt(text: string): bigint { - const match = /^\s*(\d+)/.exec(text); - if (!match) { - return 0n; - } - return BigInt(match[1]); -} - function isNumeric(a: SqliteValue): a is number | bigint { return typeof a == 'number' || typeof a == 'bigint'; } diff --git a/packages/sync-rules/src/streams/filter.ts b/packages/sync-rules/src/streams/filter.ts index 4ef22eab6..e211f377e 100644 --- a/packages/sync-rules/src/streams/filter.ts +++ b/packages/sync-rules/src/streams/filter.ts @@ -540,10 +540,11 @@ export class SubqueryParameterLookupSource implements ParameterIndexLookupCreato private streamName: string ) {} - public get defaultLookupScope() { + public get defaultLookupScope(): ParameterLookupScope { return { lookupName: this.streamName, - queryId: this.defaultQueryId + queryId: this.defaultQueryId, + source: this }; } diff --git a/packages/sync-rules/src/streams/variant.ts b/packages/sync-rules/src/streams/variant.ts index e6cf1c2ca..449e09b02 100644 --- a/packages/sync-rules/src/streams/variant.ts +++ b/packages/sync-rules/src/streams/variant.ts @@ -2,11 +2,17 @@ import { BucketInclusionReason, ResolvedBucket } from '../BucketDescription.js'; import { BucketParameterQuerier, PendingQueriers } from '../BucketParameterQuerier.js'; import { BucketDataSource, BucketParameterQuerierSource, ParameterIndexLookupCreator } from '../BucketSource.js'; import { BucketDataScope } from '../HydrationState.js'; -import { CreateSourceParams, GetQuerierOptions, RequestedStream, ScopedParameterLookup } from '../index.js'; +import { + CreateSourceParams, + GetQuerierOptions, + RequestedStream, + resolvedBucket, + ScopedParameterLookup +} from '../index.js'; import { RequestParameters, SqliteJsonValue, TableRow } from '../types.js'; -import { buildBucketName, isJsonValue, JSONBucketNameSerialize } from '../utils.js'; +import { bucketDescription, isJsonValue, JSONBucketNameSerialize } from '../utils.js'; import { BucketParameter, SubqueryEvaluator } from './parameter.js'; -import { SyncStream, SyncStreamDataSource } from './stream.js'; +import { SyncStream } from './stream.js'; import { cartesianProduct } from './utils.js'; /** @@ -291,12 +297,8 @@ export class StreamVariant { reason: BucketInclusionReason, bucketScope: BucketDataScope ): ResolvedBucket { - return { - definition: stream.name, - inclusion_reasons: [reason], - bucket: buildBucketName(bucketScope, this.serializeBucketParameters(instantiation)), - priority: stream.priority - }; + const bucketInfo = bucketDescription(bucketScope, this.serializeBucketParameters(instantiation), stream.priority); + return resolvedBucket(bucketInfo, { definition: stream.name, inclusion_reasons: [reason] }); } createParameterQuerierSource( diff --git a/packages/sync-rules/src/sync_plan/engine/javascript.ts b/packages/sync-rules/src/sync_plan/engine/javascript.ts index fe8f13dd8..607b12769 100644 --- a/packages/sync-rules/src/sync_plan/engine/javascript.ts +++ b/packages/sync-rules/src/sync_plan/engine/javascript.ts @@ -1,4 +1,5 @@ import { + cast, compare, CompatibilityContext, ExpressionType, @@ -8,7 +9,7 @@ import { sqliteBool, sqliteNot } from '../../index.js'; -import { cast, evaluateOperator, SqlFunction } from '../../sql_functions.js'; +import { evaluateOperator, SqlFunction } from '../../sql_functions.js'; import { cartesianProduct } from '../../streams/utils.js'; import { generateTableValuedFunctions } from '../../TableValuedFunctions.js'; import { SqliteRow, SqliteValue } from '../../types.js'; diff --git a/packages/sync-rules/src/sync_plan/evaluator/bucket_data_source.ts b/packages/sync-rules/src/sync_plan/evaluator/bucket_data_source.ts index 924aaeba5..5ca3573c3 100644 --- a/packages/sync-rules/src/sync_plan/evaluator/bucket_data_source.ts +++ b/packages/sync-rules/src/sync_plan/evaluator/bucket_data_source.ts @@ -9,7 +9,8 @@ import { UnscopedEvaluatedRow, UnscopedEvaluationResult } from '../../types.js'; -import { filterJsonRow, idFromData, isJsonValue, isValidParameterValue, JSONBucketNameSerialize } from '../../utils.js'; +import { idFromData } from '../../cast.js'; +import { filterJsonRow, isJsonValue, isValidParameterValue, JSONBucketNameSerialize } from '../../utils.js'; import { SqlExpression } from '../expression.js'; import { ExpressionToSqlite } from '../expression_to_sql.js'; import * as plan from '../plan.js'; diff --git a/packages/sync-rules/src/sync_plan/evaluator/bucket_source.ts b/packages/sync-rules/src/sync_plan/evaluator/bucket_source.ts index 68dbcef8d..f11568bfe 100644 --- a/packages/sync-rules/src/sync_plan/evaluator/bucket_source.ts +++ b/packages/sync-rules/src/sync_plan/evaluator/bucket_source.ts @@ -1,3 +1,5 @@ +import { BucketInclusionReason, ResolvedBucket } from '../../BucketDescription.js'; +import { PendingQueriers } from '../../BucketParameterQuerier.js'; import { BucketDataSource, BucketSource, @@ -6,16 +8,14 @@ import { HydratedBucketSource, ParameterIndexLookupCreator } from '../../BucketSource.js'; -import { StreamEvaluationContext } from './index.js'; -import * as plan from '../plan.js'; +import { RequestedStream } from '../../SqlSyncRules.js'; +import { RequestParameters, SqliteParameterValue } from '../../types.js'; +import { bucketDescription, JSONBucketNameSerialize, resolvedBucket } from '../../utils.js'; import { mapExternalDataToInstantiation, ScalarExpressionEngine } from '../engine/scalar_expression_engine.js'; import { SqlExpression } from '../expression.js'; -import { RequestParameters, SqliteParameterValue } from '../../types.js'; +import * as plan from '../plan.js'; +import { StreamEvaluationContext } from './index.js'; import { parametersForRequest, RequestParameterEvaluators } from './parameter_evaluator.js'; -import { PendingQueriers } from '../../BucketParameterQuerier.js'; -import { RequestedStream } from '../../SqlSyncRules.js'; -import { BucketInclusionReason, ResolvedBucket } from '../../BucketDescription.js'; -import { buildBucketName, JSONBucketNameSerialize } from '../../utils.js'; export interface StreamInput extends StreamEvaluationContext { preparedBuckets: Map; @@ -130,12 +130,15 @@ class PreparedQuerier { const bucketScope = hydration.hydrationState.getBucketSourceScope(this.dataSource); const parametersToBucket = (instantiation: SqliteParameterValue[]): ResolvedBucket => { - return { + const desc = bucketDescription( + bucketScope, + JSONBucketNameSerialize.stringify(instantiation), + this.stream.priority + ); + return resolvedBucket(desc, { definition: this.stream.name, - inclusion_reasons: [reason], - bucket: buildBucketName(bucketScope, JSONBucketNameSerialize.stringify(instantiation)), - priority: this.stream.priority - }; + inclusion_reasons: [reason] + }); }; // Do we need parameter lookups to resolve parameters? diff --git a/packages/sync-rules/src/sync_plan/evaluator/parameter_index_lookup_creator.ts b/packages/sync-rules/src/sync_plan/evaluator/parameter_index_lookup_creator.ts index 66785f640..4d00dc389 100644 --- a/packages/sync-rules/src/sync_plan/evaluator/parameter_index_lookup_creator.ts +++ b/packages/sync-rules/src/sync_plan/evaluator/parameter_index_lookup_creator.ts @@ -23,7 +23,10 @@ export class PreparedParameterIndexLookupCreator implements ParameterIndexLookup private readonly source: plan.StreamParameterIndexLookupCreator, { engine, defaultSchema }: StreamEvaluationContext ) { - this.defaultLookupScope = source.defaultLookupScope; + this.defaultLookupScope = { + ...source.defaultLookupScope, + source: this + }; const translationHelper = new TableProcessorToSqlHelper(source); const expressions = source.outputs.map((o) => translationHelper.mapper.transform(o)); diff --git a/packages/sync-rules/src/sync_plan/plan.ts b/packages/sync-rules/src/sync_plan/plan.ts index 17edfcf36..53b411f82 100644 --- a/packages/sync-rules/src/sync_plan/plan.ts +++ b/packages/sync-rules/src/sync_plan/plan.ts @@ -161,7 +161,7 @@ export interface StreamBucketDataSource { */ export interface StreamParameterIndexLookupCreator extends TableProcessor { hashCode: number; - defaultLookupScope: ParameterLookupScope; + defaultLookupScope: Omit; /** * Outputs to persist in the lookup. diff --git a/packages/sync-rules/src/sync_plan/serialize.ts b/packages/sync-rules/src/sync_plan/serialize.ts index 3f21be556..2f7887a8c 100644 --- a/packages/sync-rules/src/sync_plan/serialize.ts +++ b/packages/sync-rules/src/sync_plan/serialize.ts @@ -378,7 +378,7 @@ interface SerializedDataSource { interface SerializedParameterIndexLookupCreator { table: SerializedTablePattern; hash: number; - lookupScope: ParameterLookupScope; + lookupScope: Omit; output: SqlExpression[]; filters: SqlExpression[]; tableValuedFunctions: TableProcessorTableValuedFunction[]; diff --git a/packages/sync-rules/src/types.ts b/packages/sync-rules/src/types.ts index 6f55647b3..1b2076d42 100644 --- a/packages/sync-rules/src/types.ts +++ b/packages/sync-rules/src/types.ts @@ -7,6 +7,7 @@ import { RequestFunctionCall } from './request_functions.js'; import { SourceTableInterface } from './SourceTableInterface.js'; import { SyncRulesOptions } from './SqlSyncRules.js'; import { TablePattern } from './TablePattern.js'; +import { BucketDataSource } from './BucketSource.js'; import { CustomSqliteValue } from './types/custom_sqlite_value.js'; import { jsonValueToSqlite, toSyncRulesParameters } from './utils.js'; @@ -58,6 +59,9 @@ export interface EvaluatedRow { /** Must be JSON-serializable. */ data: SqliteJsonRow; + + /** Source for the evaluated row. */ + source: BucketDataSource; } /** diff --git a/packages/sync-rules/src/utils.ts b/packages/sync-rules/src/utils.ts index 216f56847..51a9c4e88 100644 --- a/packages/sync-rules/src/utils.ts +++ b/packages/sync-rules/src/utils.ts @@ -1,5 +1,7 @@ import { JSONBig, JsonContainer, Replacer, stringifyRaw } from '@powersync/service-jsonbig'; import { SelectFromStatement, Statement } from 'pgsql-ast-parser'; +import { BucketDescription, BucketInclusionReason, BucketPriority, ResolvedBucket } from './BucketDescription.js'; +import { BucketDataSource } from './BucketSource.js'; import { CompatibilityContext } from './compatibility.js'; import { SyncRuleProcessingError as SyncRulesProcessingError } from './errors.js'; import { BucketDataScope } from './HydrationState.js'; @@ -16,12 +18,80 @@ import { SqliteValue } from './types.js'; import { CustomArray, CustomObject, CustomSqliteValue } from './types/custom_sqlite_value.js'; -import { castAsText } from './sql_functions.js'; export function isSelectStatement(q: Statement): q is SelectFromStatement { return q.type == 'select'; } +export function bucketDescription( + scope: BucketDataScope, + serializedParameters: string, + priority: BucketPriority +): BucketDescription { + const info = { bucket: scope.bucketPrefix + serializedParameters, priority }; + return withBucketSource(info, scope.source); +} + +export function resolvedBucket( + description: BucketDescription, + options: { definition: string; inclusion_reasons: BucketInclusionReason[] } +): ResolvedBucket { + const result = { + ...description, + ...options + }; + return withBucketSource(result, description.source); +} + +/** + * Resolves duplicate buckets in the given array, merging the inclusion reasons for duplicate. + * + * It's possible for duplicates to occur when a stream has multiple subscriptions, consider e.g. + * + * ``` + * sync_streams: + * assets_by_category: + * query: select * from assets where category in (request.parameters() -> 'categories') + * ``` + * + * Here, a client might subscribe once with `{"categories": [1]}` and once with `{"categories": [1, 2]}`. Since each + * subscription is evaluated independently, this would lead to three buckets, with a duplicate `assets_by_category[1]` + * bucket. + */ +export function mergeBuckets(buckets: ResolvedBucket[]): ResolvedBucket[] { + const byBucketId: Record = {}; + + for (const bucket of buckets) { + if (Object.hasOwn(byBucketId, bucket.bucket)) { + byBucketId[bucket.bucket].inclusion_reasons.push(...bucket.inclusion_reasons); + } else { + // Clone so that we can modify the merged value without affecting the input value + byBucketId[bucket.bucket] = cloneResolvedBucket(bucket); + } + } + + return Object.values(byBucketId); +} + +function cloneResolvedBucket(bucket: ResolvedBucket) { + let clone = structuredClone(bucket); + // The structured clone does not include the non-enumerable source - set it directly. + return withBucketSource(clone, bucket.source); +} + +export function withBucketSource( + value: T, + source: BucketDataSource +): T & { source: BucketDataSource } { + Object.defineProperty(value, 'source', { + value: source, + // This is important. If the property is enumerable, it may end up in JSON output to the client, + // and will pollute tests. + enumerable: false + }); + return value as T & { source: BucketDataSource }; +} + export function buildBucketName(scope: BucketDataScope, serializedParameters: string): string { return scope.bucketPrefix + serializedParameters; } @@ -239,20 +309,3 @@ export function normalizeParameterValue(value: SqliteJsonValue): SqliteJsonValue } return value; } - -/** - * Extracts and normalizes the ID column from a row. - */ -export function idFromData(data: SqliteJsonRow): string { - let id = data.id; - if (typeof id != 'string') { - // While an explicit cast would be better, this covers against very common - // issues when initially testing out sync, for example when the id column is an - // auto-incrementing integer. - // If there is no id column, we use a blank id. This will result in the user syncing - // a single arbitrary row for this table - better than just not being able to sync - // anything. - id = castAsText(id) ?? ''; - } - return id; -} diff --git a/packages/sync-rules/test/src/parameter_queries.test.ts b/packages/sync-rules/test/src/parameter_queries.test.ts index e1e703bce..3bc74bfca 100644 --- a/packages/sync-rules/test/src/parameter_queries.test.ts +++ b/packages/sync-rules/test/src/parameter_queries.test.ts @@ -2,8 +2,10 @@ import { beforeEach, describe, expect, test } from 'vitest'; import { HydrationState } from '../../src/HydrationState.js'; import { BucketParameterQuerier, + BucketDataScope, GetQuerierOptions, mergeParameterIndexLookupCreators, + ParameterLookupScope, QuerierError, RequestParameters, ScopedParameterLookup, @@ -12,7 +14,15 @@ import { UnscopedParameterLookup } from '../../src/index.js'; import { StaticSqlParameterQuery } from '../../src/StaticSqlParameterQuery.js'; -import { BASIC_SCHEMA, EMPTY_DATA_SOURCE, findQuerierLookups, PARSE_OPTIONS, requestParameters } from './util.js'; +import { + BASIC_SCHEMA, + bucketDataScope, + EMPTY_DATA_SOURCE, + findQuerierLookups, + lookupScope, + PARSE_OPTIONS, + requestParameters +} from './util.js'; describe('parameter queries', () => { const table = (name: string): SourceTableInterface => ({ @@ -123,15 +133,19 @@ describe('parameter queries', () => { // We _do_ need to care about the bucket string representation. expect( - query.resolveBucketDescriptions([{ int1: 314, float1: 3.14, float2: 314 }], requestParameters({}), { - bucketPrefix: 'mybucket' - }) + query.resolveBucketDescriptions( + [{ int1: 314, float1: 3.14, float2: 314 }], + requestParameters({}), + bucketDataScope('mybucket') + ) ).toEqual([{ bucket: 'mybucket[314,3.14,314]', priority: 3 }]); expect( - query.resolveBucketDescriptions([{ int1: 314n, float1: 3.14, float2: 314 }], requestParameters({}), { - bucketPrefix: 'mybucket' - }) + query.resolveBucketDescriptions( + [{ int1: 314n, float1: 3.14, float2: 314 }], + requestParameters({}), + bucketDataScope('mybucket') + ) ).toEqual([{ bucket: 'mybucket[314,3.14,314]', priority: 3 }]); }); @@ -494,7 +508,7 @@ describe('parameter queries', () => { query.resolveBucketDescriptions( [{ user_id: 'user1' }], requestParameters({ sub: 'user1', parameters: { is_admin: true } }), - { bucketPrefix: 'mybucket' } + bucketDataScope('mybucket') ) ).toEqual([{ bucket: 'mybucket["user1",1]', priority: 3 }]); }); @@ -872,13 +886,14 @@ describe('parameter queries', () => { describe('custom hydrationState', function () { const hydrationState: HydrationState = { - getBucketSourceScope(source) { - return { bucketPrefix: `${source.uniqueName}-test` }; + getBucketSourceScope(source): BucketDataScope { + return { bucketPrefix: `${source.uniqueName}-test`, source }; }, - getParameterIndexLookupScope(source) { + getParameterIndexLookupScope(source): ParameterLookupScope { return { lookupName: `${source.defaultLookupScope.lookupName}.test`, - queryId: `${source.defaultLookupScope.queryId}.test` + queryId: `${source.defaultLookupScope.queryId}.test`, + source }; } }; @@ -906,13 +921,11 @@ describe('parameter queries', () => { }); expect(result).toEqual([ { - lookup: ScopedParameterLookup.direct({ lookupName: 'mybucket.test', queryId: 'myquery.test' }, ['test-user']), + lookup: ScopedParameterLookup.direct(lookupScope('mybucket.test', 'myquery.test'), ['test-user']), bucketParameters: [{ group_id: 'group1' }] }, { - lookup: ScopedParameterLookup.direct({ lookupName: 'mybucket.test', queryId: 'myquery.test' }, [ - 'other-user' - ]), + lookup: ScopedParameterLookup.direct(lookupScope('mybucket.test', 'myquery.test'), ['other-user']), bucketParameters: [{ group_id: 'group1' }] } ]); @@ -944,7 +957,7 @@ describe('parameter queries', () => { const querier = queriers[0]; expect(querier.hasDynamicBuckets).toBeTruthy(); expect(await findQuerierLookups(querier)).toEqual([ - ScopedParameterLookup.direct({ lookupName: 'mybucket.test', queryId: 'myquery.test' }, ['test-user']) + ScopedParameterLookup.direct(lookupScope('mybucket.test', 'myquery.test'), ['test-user']) ]); }); }); diff --git a/packages/sync-rules/test/src/static_parameter_queries.test.ts b/packages/sync-rules/test/src/static_parameter_queries.test.ts index 6b6c66dcf..b91bf35e9 100644 --- a/packages/sync-rules/test/src/static_parameter_queries.test.ts +++ b/packages/sync-rules/test/src/static_parameter_queries.test.ts @@ -1,12 +1,13 @@ import { describe, expect, test } from 'vitest'; -import { BucketDataScope, HydrationState } from '../../src/HydrationState.js'; +import { BucketDataScope, HydrationState, ParameterLookupScope } from '../../src/HydrationState.js'; import { BucketParameterQuerier, GetQuerierOptions, QuerierError, SqlParameterQuery } from '../../src/index.js'; import { StaticSqlParameterQuery } from '../../src/StaticSqlParameterQuery.js'; -import { EMPTY_DATA_SOURCE, PARSE_OPTIONS, requestParameters } from './util.js'; +import { bucketDataScope, EMPTY_DATA_SOURCE, PARSE_OPTIONS, requestParameters } from './util.js'; describe('static parameter queries', () => { const MYBUCKET_SCOPE: BucketDataScope = { - bucketPrefix: 'mybucket' + bucketPrefix: 'mybucket', + source: EMPTY_DATA_SOURCE }; test('basic query', function () { @@ -37,9 +38,7 @@ describe('static parameter queries', () => { expect(query.errors).toEqual([]); expect(query.bucketParameters!).toEqual(['user_id']); expect( - query.getStaticBucketDescriptions(requestParameters({ sub: 'user1' }), { - bucketPrefix: '1#mybucket' - }) + query.getStaticBucketDescriptions(requestParameters({ sub: 'user1' }), bucketDataScope('1#mybucket')) ).toEqual([{ bucket: '1#mybucket["user1"]', priority: 3 }]); }); @@ -478,13 +477,14 @@ describe('static parameter queries', () => { expect(query.errors).toEqual([]); const hydrationState: HydrationState = { - getBucketSourceScope(source) { - return { bucketPrefix: `${source.uniqueName}-test` }; + getBucketSourceScope(source): BucketDataScope { + return { bucketPrefix: `${source.uniqueName}-test`, source }; }, - getParameterIndexLookupScope(source) { + getParameterIndexLookupScope(source): ParameterLookupScope { return { lookupName: `${source.defaultLookupScope.lookupName}.test`, - queryId: `${source.defaultLookupScope.queryId}.test` + queryId: `${source.defaultLookupScope.queryId}.test`, + source }; } }; diff --git a/packages/sync-rules/test/src/streams.test.ts b/packages/sync-rules/test/src/streams.test.ts index 6f5709b12..137f9344d 100644 --- a/packages/sync-rules/test/src/streams.test.ts +++ b/packages/sync-rules/test/src/streams.test.ts @@ -3,6 +3,7 @@ import { describe, expect, test } from 'vitest'; import { HydrationState, ParameterLookupScope, versionedHydrationState } from '../../src/HydrationState.js'; import { BucketParameterQuerier, + BucketDataScope, CompatibilityContext, CompatibilityEdition, CreateSourceParams, @@ -14,7 +15,6 @@ import { mergeBucketParameterQueriers, UnscopedParameterLookup, QuerierError, - RequestParameters, SourceTableInterface, SqliteJsonRow, SqliteRow, @@ -24,17 +24,11 @@ import { syncStreamFromSql, ScopedParameterLookup } from '../../src/index.js'; -import { normalizeQuerierOptions, PARSE_OPTIONS, requestParameters, TestSourceTable } from './util.js'; +import { lookupScope, normalizeQuerierOptions, PARSE_OPTIONS, requestParameters, TestSourceTable } from './util.js'; describe('streams', () => { - const STREAM_0: ParameterLookupScope = { - lookupName: 'stream', - queryId: '0' - }; - const STREAM_1: ParameterLookupScope = { - lookupName: 'stream', - queryId: '1' - }; + const STREAM_0: ParameterLookupScope = lookupScope('stream', '0'); + const STREAM_1: ParameterLookupScope = lookupScope('stream', '1'); test('refuses edition: 1', () => { expect(() => @@ -104,6 +98,49 @@ describe('streams', () => { ]); }); + test('tracks source metadata on stream APIs', async () => { + const desc = parseStream( + "SELECT * FROM comments WHERE issue_id = subscription.parameter('issue') OR issue_id IN (SELECT id FROM issues WHERE owner_id = auth.user_id())" + ); + const source = debugHydratedMergedSource(desc, hydrationParams); + + const rowResults = source.evaluateRow({ sourceTable: COMMENTS, record: { id: 'c1', issue_id: 'i1' } }); + expect(rowResults).toHaveLength(2); + + const rowSourceByBucket = new Map(); + for (const row of rowResults) { + if ('error' in row) { + throw new Error(`Unexpected error evaluating stream row: ${row.error}`); + } + rowSourceByBucket.set(row.bucket, row.source); + } + expect(rowSourceByBucket.get('1#stream|0["i1"]')).toBe(desc.dataSources[0]); + expect(rowSourceByBucket.get('1#stream|1["i1"]')).toBe(desc.dataSources[1]); + + const parameterResults = source.evaluateParameterRow(ISSUES, { id: 'i1', owner_id: 'u1' }); + expect(parameterResults).toHaveLength(1); + if ('error' in parameterResults[0]) { + throw new Error(`Unexpected error evaluating stream parameters: ${parameterResults[0].error}`); + } + expect(parameterResults[0].lookup.source).toBe(desc.parameterIndexLookupCreators[0]); + + const { querier, errors } = await createQueriers(desc, { + tokenPayload: { sub: 'u1' }, + parameters: { issue: 'i1' } + }); + expect(errors).toHaveLength(0); + expect(querier.staticBuckets).toHaveLength(1); + expect(querier.staticBuckets[0].source).toBe(desc.dataSources[0]); + + const dynamicBuckets = await querier.queryDynamicBucketDescriptions({ + async getParameterSets() { + return [{ result: 'i1' }]; + } + }); + expect(dynamicBuckets).toHaveLength(1); + expect(dynamicBuckets[0].source).toBe(desc.dataSources[1]); + }); + describe('or', () => { test('parameter match or request condition', async () => { const desc = parseStream("SELECT * FROM issues WHERE owner_id = auth.user_id() OR auth.parameter('is_admin')"); @@ -759,9 +796,7 @@ describe('streams', () => { tokenPayload: { sub: 'id' }, parameters: {}, getParameterSets(lookups) { - expect(lookups).toStrictEqual([ - ScopedParameterLookup.direct({ lookupName: 'account_member', queryId: '0' }, ['id']) - ]); + expect(lookups).toStrictEqual([ScopedParameterLookup.direct(lookupScope('account_member', '0'), ['id'])]); return [{ result: 'account_id' }]; } }) @@ -972,13 +1007,14 @@ WHERE `); const hydrationState: HydrationState = { - getBucketSourceScope(source) { - return { bucketPrefix: `${source.uniqueName}.test` }; + getBucketSourceScope(source): BucketDataScope { + return { bucketPrefix: `${source.uniqueName}.test`, source }; }, - getParameterIndexLookupScope(source) { + getParameterIndexLookupScope(source): ParameterLookupScope { return { lookupName: `${source.defaultLookupScope.lookupName}.test`, - queryId: `${source.defaultLookupScope.queryId}.test` + queryId: `${source.defaultLookupScope.queryId}.test`, + source }; } }; @@ -997,7 +1033,7 @@ WHERE }) ).toStrictEqual([ { - lookup: ScopedParameterLookup.direct({ lookupName: 'stream.test', queryId: '0.test' }, ['u1']), + lookup: ScopedParameterLookup.direct(lookupScope('stream.test', '0.test'), ['u1']), bucketParameters: [ { result: 'i1' @@ -1006,7 +1042,7 @@ WHERE }, { - lookup: ScopedParameterLookup.direct({ lookupName: 'stream.test', queryId: '1.test' }, ['myname']), + lookup: ScopedParameterLookup.direct(lookupScope('stream.test', '1.test'), ['myname']), bucketParameters: [ { result: 'i1' @@ -1022,7 +1058,7 @@ WHERE }) ).toStrictEqual([ { - lookup: ScopedParameterLookup.direct({ lookupName: 'stream.test', queryId: '0.test' }, ['u1']), + lookup: ScopedParameterLookup.direct(lookupScope('stream.test', '0.test'), ['u1']), bucketParameters: [ { result: 'i1' 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 42ca08aa4..e2fdde80d 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 @@ -8,7 +8,7 @@ import { SqliteRow, SqliteValue } from '../../../../src/index.js'; -import { requestParameters, TestSourceTable } from '../../util.js'; +import { lookupScope, requestParameters, TestSourceTable } from '../../util.js'; describe('evaluating rows', () => { syncTest('emits rows', ({ sync }) => { @@ -222,7 +222,7 @@ streams: expect(desc.evaluateParameterRow(ISSUES, { id: 'issue_id', owner_id: 'user1', name: 'name' })).toStrictEqual([ { - lookup: ScopedParameterLookup.direct({ lookupName: 'lookup', queryId: '0' }, ['user1']), + lookup: ScopedParameterLookup.direct(lookupScope('lookup', '0'), ['user1']), bucketParameters: [ { '0': 'issue_id' @@ -274,6 +274,51 @@ streams: }); describe('querier', () => { + syncTest('tracks source metadata on stream APIs', async ({ sync }) => { + const desc = sync.prepareSyncStreams(` +config: + edition: 3 + +streams: + stream: + accept_potentially_dangerous_queries: true + queries: + - SELECT * FROM comments WHERE issue_id = subscription.parameter('issue') + - SELECT * FROM comments WHERE issue_id IN (SELECT id FROM issues WHERE owner_id = auth.user_id()) +`); + const streamSource = desc.definition.bucketSources[0]; + expect(streamSource.dataSources).toHaveLength(2); + + const rowResults = desc.evaluateRow({ sourceTable: COMMENTS, record: { id: 'c1', issue_id: 'i1' } }); + expect(rowResults).toHaveLength(1); + expect(rowResults[0].bucket).toBe('stream|0["i1"]'); + expect(rowResults[0].source).toBe(streamSource.dataSources[0]); + + expect(desc.definition.bucketParameterLookupSources).toHaveLength(1); + const parameterResults = desc.evaluateParameterRow(ISSUES, { id: 'i1', owner_id: 'u1' }); + expect(parameterResults).toHaveLength(1); + expect(parameterResults[0].lookup.source).toBe(desc.definition.bucketParameterLookupSources[0]); + + const { querier, errors } = desc.getBucketParameterQuerier({ + globalParameters: requestParameters({ sub: 'u1' }), + hasDefaultStreams: false, + streams: { + stream: [{ opaque_id: 0, parameters: { issue: 'i1' } }] + } + }); + expect(errors).toHaveLength(0); + expect(querier.staticBuckets).toHaveLength(1); + expect(querier.staticBuckets[0].source).toBe(streamSource.dataSources[0]); + + const dynamicBuckets = await querier.queryDynamicBucketDescriptions({ + async getParameterSets() { + return [{ '0': 'i1' }]; + } + }); + expect(dynamicBuckets).toHaveLength(1); + expect(dynamicBuckets[0].source).toBe(streamSource.dataSources[1]); + }); + syncTest('static', ({ sync }) => { const desc = sync.prepareSyncStreams(` config: @@ -349,28 +394,12 @@ streams: if (call == 0) { // First call. Lookup from users.id => users.name call++; - expect(lookups).toStrictEqual([ - ScopedParameterLookup.direct( - { - lookupName: 'lookup', - queryId: '0' - }, - ['user'] - ) - ]); + expect(lookups).toStrictEqual([ScopedParameterLookup.direct(lookupScope('lookup', '0'), ['user'])]); return [{ '0': 'name' }]; } else if (call == 1) { // Second call. Lookup from issues.owned_by => issues.id call++; - expect(lookups).toStrictEqual([ - ScopedParameterLookup.direct( - { - lookupName: 'lookup', - queryId: '1' - }, - ['name'] - ) - ]); + expect(lookups).toStrictEqual([ScopedParameterLookup.direct(lookupScope('lookup', '1'), ['name'])]); return [{ '0': 'issue' }]; } diff --git a/packages/sync-rules/test/src/sync_plan/evaluator/table_valued.test.ts b/packages/sync-rules/test/src/sync_plan/evaluator/table_valued.test.ts index 35e827d77..378f78ff0 100644 --- a/packages/sync-rules/test/src/sync_plan/evaluator/table_valued.test.ts +++ b/packages/sync-rules/test/src/sync_plan/evaluator/table_valued.test.ts @@ -1,6 +1,6 @@ import { describe, expect } from 'vitest'; import { syncTest } from './utils.js'; -import { requestParameters, TestSourceTable } from '../../util.js'; +import { lookupScope, requestParameters, TestSourceTable } from '../../util.js'; import { deserializeSyncPlan, ImplicitSchemaTablePattern, @@ -55,7 +55,7 @@ streams: desc.evaluateParameterRow(conversations, { id: 'chat', members: JSON.stringify(['user', 'another']) }) ).toStrictEqual([ { - lookup: ScopedParameterLookup.direct({ lookupName: 'lookup', queryId: '0' }, ['chat']), + lookup: ScopedParameterLookup.direct(lookupScope('lookup', '0'), ['chat']), bucketParameters: [ { '0': 'user' @@ -63,7 +63,7 @@ streams: ] }, { - lookup: ScopedParameterLookup.direct({ lookupName: 'lookup', queryId: '0' }, ['chat']), + lookup: ScopedParameterLookup.direct(lookupScope('lookup', '0'), ['chat']), bucketParameters: [ { '0': 'another' @@ -87,15 +87,7 @@ streams: const buckets = await querier.queryDynamicBucketDescriptions({ getParameterSets: async function (lookups: ScopedParameterLookup[]): Promise { - expect(lookups).toStrictEqual([ - ScopedParameterLookup.direct( - { - lookupName: 'lookup', - queryId: '0' - }, - ['chat'] - ) - ]); + expect(lookups).toStrictEqual([ScopedParameterLookup.direct(lookupScope('lookup', '0'), ['chat'])]); return [{ '0': 'user' }, { '0': 'another' }]; } diff --git a/packages/sync-rules/test/src/sync_rules.test.ts b/packages/sync-rules/test/src/sync_rules.test.ts index 44d3c2066..a7c3c3a7a 100644 --- a/packages/sync-rules/test/src/sync_rules.test.ts +++ b/packages/sync-rules/test/src/sync_rules.test.ts @@ -1,18 +1,24 @@ import { describe, expect, test } from 'vitest'; import { CreateSourceParams, ScopedParameterLookup, SqlSyncRules } from '../../src/index.js'; -import { DEFAULT_HYDRATION_STATE, HydrationState } from '../../src/HydrationState.js'; +import { + BucketDataScope, + DEFAULT_HYDRATION_STATE, + HydrationState, + ParameterLookupScope +} from '../../src/HydrationState.js'; import { SqlBucketDescriptor } from '../../src/SqlBucketDescriptor.js'; import { StaticSqlParameterQuery } from '../../src/StaticSqlParameterQuery.js'; import { ASSETS, BASIC_SCHEMA, - PARSE_OPTIONS, - TestSourceTable, - USERS, findQuerierLookups, + lookupScope, normalizeQuerierOptions, - requestParameters + PARSE_OPTIONS, + requestParameters, + TestSourceTable, + USERS } from './util.js'; describe('sync rules', () => { @@ -63,6 +69,38 @@ bucket_definitions: }); }); + test('tracks source metadata on rows, lookups and bucket descriptions', () => { + const { config: rules } = SqlSyncRules.fromYaml( + ` +bucket_definitions: + mybucket: + parameters: + - SELECT token_parameters.user_id as user_id + - SELECT users.id as user_id FROM users WHERE users.id = token_parameters.user_id + data: + - SELECT id FROM assets WHERE assets.user_id = bucket.user_id + `, + PARSE_OPTIONS + ); + const hydrated = rules.hydrate(hydrationParams); + + const staticBuckets = hydrated.getBucketParameterQuerier(normalizeQuerierOptions({ sub: 'user1' })).querier + .staticBuckets; + expect(staticBuckets).toHaveLength(1); + expect(staticBuckets[0].source).toBe(rules.bucketDataSources[0]); + + const dataResults = hydrated.evaluateRow({ + sourceTable: ASSETS, + record: { id: 'asset1', user_id: 'user1' } + }); + expect(dataResults).toHaveLength(1); + expect(dataResults[0].source).toBe(rules.bucketDataSources[0]); + + const parameterResults = hydrated.evaluateParameterRow(USERS, { id: 'user1' }); + expect(parameterResults).toHaveLength(1); + expect(parameterResults[0].lookup.source).toBe(rules.bucketParameterLookupSources[0]); + }); + test('parse global sync rules with filter', () => { const { config: rules } = SqlSyncRules.fromYaml( ` @@ -114,7 +152,7 @@ bucket_definitions: expect(hydrated.evaluateParameterRow(USERS, { id: 'user1', is_admin: 1 })).toEqual([ { bucketParameters: [{}], - lookup: ScopedParameterLookup.direct({ lookupName: 'mybucket', queryId: '1' }, ['user1']) + lookup: ScopedParameterLookup.direct(lookupScope('mybucket', '1'), ['user1']) } ]); expect(hydrated.evaluateParameterRow(USERS, { id: 'user1', is_admin: 0 })).toEqual([]); @@ -183,13 +221,14 @@ bucket_definitions: PARSE_OPTIONS ); const hydrationState: HydrationState = { - getBucketSourceScope(source) { - return { bucketPrefix: `${source.uniqueName}-test` }; + getBucketSourceScope(source): BucketDataScope { + return { bucketPrefix: `${source.uniqueName}-test`, source }; }, - getParameterIndexLookupScope(source) { + getParameterIndexLookupScope(source): ParameterLookupScope { return { lookupName: `${source.defaultLookupScope.lookupName}.test`, - queryId: `${source.defaultLookupScope.queryId}.test` + queryId: `${source.defaultLookupScope.queryId}.test`, + source }; } }; @@ -207,13 +246,13 @@ bucket_definitions: } ]); expect(await findQuerierLookups(querier)).toEqual([ - ScopedParameterLookup.direct({ lookupName: 'mybucket.test', queryId: '2.test' }, ['user1']) + ScopedParameterLookup.direct(lookupScope('mybucket.test', '2.test'), ['user1']) ]); expect(hydrated.evaluateParameterRow(USERS, { id: 'user1', is_admin: 1 })).toEqual([ { bucketParameters: [{ user_id: 'user1' }], - lookup: ScopedParameterLookup.direct({ lookupName: 'mybucket.test', queryId: '2.test' }, ['user1']) + lookup: ScopedParameterLookup.direct(lookupScope('mybucket.test', '2.test'), ['user1']) } ]); @@ -1044,7 +1083,7 @@ bucket_definitions: }); expect(await findQuerierLookups(hydratedQuerier)).toEqual([ - ScopedParameterLookup.direct({ lookupName: 'admin_only', queryId: '1' }, [1]) + ScopedParameterLookup.direct(lookupScope('admin_only', '1'), [1]) ]); }); diff --git a/packages/sync-rules/test/src/table_valued_function_queries.test.ts b/packages/sync-rules/test/src/table_valued_function_queries.test.ts index 411e1d701..f24c9902a 100644 --- a/packages/sync-rules/test/src/table_valued_function_queries.test.ts +++ b/packages/sync-rules/test/src/table_valued_function_queries.test.ts @@ -8,7 +8,7 @@ import { SqlParameterQuery } from '../../src/index.js'; import { StaticSqlParameterQuery } from '../../src/StaticSqlParameterQuery.js'; -import { EMPTY_DATA_SOURCE, PARSE_OPTIONS, requestParameters } from './util.js'; +import { bucketDataScope, EMPTY_DATA_SOURCE, PARSE_OPTIONS, requestParameters } from './util.js'; describe('table-valued function queries', () => { const emptyPayload: RequestJwtPayload = { userIdJson: '', parsedPayload: {} }; @@ -29,9 +29,7 @@ describe('table-valued function queries', () => { expect(query.bucketParameters).toEqual(['v']); expect( - query.getStaticBucketDescriptions(requestParameters({}, { array: [1, 2, 3, null] }), { - bucketPrefix: 'mybucket' - }) + query.getStaticBucketDescriptions(requestParameters({}, { array: [1, 2, 3, null] }), bucketDataScope('mybucket')) ).toEqual([ { bucket: 'mybucket[1]', priority: 3 }, { bucket: 'mybucket[2]', priority: 3 }, @@ -60,9 +58,7 @@ describe('table-valued function queries', () => { expect(query.bucketParameters).toEqual(['v']); expect( - query.getStaticBucketDescriptions(requestParameters({}, { array: [1, 2, 3, null] }), { - bucketPrefix: 'mybucket' - }) + query.getStaticBucketDescriptions(requestParameters({}, { array: [1, 2, 3, null] }), bucketDataScope('mybucket')) ).toEqual([ { bucket: 'mybucket[1]', priority: 3 }, { bucket: 'mybucket[2]', priority: 3 }, @@ -83,11 +79,7 @@ describe('table-valued function queries', () => { expect(query.errors).toEqual([]); expect(query.bucketParameters).toEqual(['v']); - expect( - query.getStaticBucketDescriptions(requestParameters({}, {}), { - bucketPrefix: 'mybucket' - }) - ).toEqual([ + expect(query.getStaticBucketDescriptions(requestParameters({}, {}), bucketDataScope('mybucket'))).toEqual([ { bucket: 'mybucket[1]', priority: 3 }, { bucket: 'mybucket[2]', priority: 3 }, { bucket: 'mybucket[3]', priority: 3 } @@ -106,11 +98,7 @@ describe('table-valued function queries', () => { expect(query.errors).toEqual([]); expect(query.bucketParameters).toEqual(['v']); - expect( - query.getStaticBucketDescriptions(requestParameters({}, {}), { - bucketPrefix: 'mybucket' - }) - ).toEqual([]); + expect(query.getStaticBucketDescriptions(requestParameters({}, {}), bucketDataScope('mybucket'))).toEqual([]); }); test('json_each(array param not present)', function () { @@ -128,11 +116,7 @@ describe('table-valued function queries', () => { expect(query.errors).toEqual([]); expect(query.bucketParameters).toEqual(['v']); - expect( - query.getStaticBucketDescriptions(requestParameters({}, {}), { - bucketPrefix: 'mybucket' - }) - ).toEqual([]); + expect(query.getStaticBucketDescriptions(requestParameters({}, {}), bucketDataScope('mybucket'))).toEqual([]); }); test('json_each(array param not present, ifnull)', function () { @@ -150,11 +134,7 @@ describe('table-valued function queries', () => { expect(query.errors).toEqual([]); expect(query.bucketParameters).toEqual(['v']); - expect( - query.getStaticBucketDescriptions(requestParameters({}, {}), { - bucketPrefix: 'mybucket' - }) - ).toEqual([]); + expect(query.getStaticBucketDescriptions(requestParameters({}, {}), bucketDataScope('mybucket'))).toEqual([]); }); test('json_each on json_keys', function () { @@ -169,11 +149,7 @@ describe('table-valued function queries', () => { expect(query.errors).toEqual([]); expect(query.bucketParameters).toEqual(['value']); - expect( - query.getStaticBucketDescriptions(requestParameters({}, {}), { - bucketPrefix: 'mybucket' - }) - ).toEqual([ + expect(query.getStaticBucketDescriptions(requestParameters({}, {}), bucketDataScope('mybucket'))).toEqual([ { bucket: 'mybucket["a"]', priority: 3 }, { bucket: 'mybucket["b"]', priority: 3 }, { bucket: 'mybucket["c"]', priority: 3 } @@ -196,9 +172,7 @@ describe('table-valued function queries', () => { expect(query.bucketParameters).toEqual(['value']); expect( - query.getStaticBucketDescriptions(requestParameters({}, { array: [1, 2, 3] }), { - bucketPrefix: 'mybucket' - }) + query.getStaticBucketDescriptions(requestParameters({}, { array: [1, 2, 3] }), bucketDataScope('mybucket')) ).toEqual([ { bucket: 'mybucket[1]', priority: 3 }, { bucket: 'mybucket[2]', priority: 3 }, @@ -222,9 +196,7 @@ describe('table-valued function queries', () => { expect(query.bucketParameters).toEqual(['value']); expect( - query.getStaticBucketDescriptions(requestParameters({}, { array: [1, 2, 3] }), { - bucketPrefix: 'mybucket' - }) + query.getStaticBucketDescriptions(requestParameters({}, { array: [1, 2, 3] }), bucketDataScope('mybucket')) ).toEqual([ { bucket: 'mybucket[1]', priority: 3 }, { bucket: 'mybucket[2]', priority: 3 }, @@ -248,9 +220,7 @@ describe('table-valued function queries', () => { expect(query.bucketParameters).toEqual(['v']); expect( - query.getStaticBucketDescriptions(requestParameters({}, { array: [1, 2, 3] }), { - bucketPrefix: 'mybucket' - }) + query.getStaticBucketDescriptions(requestParameters({}, { array: [1, 2, 3] }), bucketDataScope('mybucket')) ).toEqual([ { bucket: 'mybucket[2]', priority: 3 }, { bucket: 'mybucket[3]', priority: 3 } @@ -284,9 +254,7 @@ describe('table-valued function queries', () => { }, {} ), - { - bucketPrefix: 'mybucket' - } + bucketDataScope('mybucket') ) ).toEqual([{ bucket: 'mybucket[1]', priority: 3 }]); }); diff --git a/packages/sync-rules/test/src/util.ts b/packages/sync-rules/test/src/util.ts index f7bb60312..fe3fdd72c 100644 --- a/packages/sync-rules/test/src/util.ts +++ b/packages/sync-rules/test/src/util.ts @@ -2,12 +2,15 @@ import { expect } from 'vitest'; import { BaseJwtPayload, BucketDataSource, + BucketDataScope, BucketParameterQuerier, ColumnDefinition, CompatibilityContext, CreateSourceParams, DEFAULT_TAG, GetQuerierOptions, + ParameterIndexLookupCreator, + ParameterLookupScope, RequestedStream, RequestJwtPayload, RequestParameters, @@ -110,6 +113,37 @@ export const EMPTY_DATA_SOURCE: BucketDataSource = { } }; +export const EMPTY_PARAMETER_LOOKUP_SOURCE: ParameterIndexLookupCreator = { + get defaultLookupScope(): ParameterLookupScope { + return { + lookupName: 'lookup', + queryId: '0', + source: EMPTY_PARAMETER_LOOKUP_SOURCE + }; + }, + getSourceTables(): Set { + return new Set(); + }, + evaluateParameterRow() { + return []; + }, + tableSyncsParameters() { + return false; + } +}; + +export function bucketDataScope(bucketPrefix: string, source: BucketDataSource = EMPTY_DATA_SOURCE): BucketDataScope { + return { bucketPrefix, source }; +} + +export function lookupScope( + lookupName: string, + queryId: string, + source: ParameterIndexLookupCreator = EMPTY_PARAMETER_LOOKUP_SOURCE +): ParameterLookupScope { + return { lookupName, queryId, source }; +} + export async function findQuerierLookups(querier: BucketParameterQuerier): Promise { expect(querier.hasDynamicBuckets).toBe(true);