Skip to content

Commit a41deb1

Browse files
committed
On-demand sync - client side implementation.
1 parent bfd2794 commit a41deb1

6 files changed

Lines changed: 2428 additions & 52 deletions

File tree

‎packages/powersync-db-collection/package.json‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -55,6 +55,7 @@
5555
"@standard-schema/spec": "^1.1.0",
5656
"@tanstack/db": "workspace:*",
5757
"@tanstack/store": "^0.8.0",
58+
"async-mutex": "^0.5.0",
5859
"debug": "^4.4.3",
5960
"p-defer": "^4.0.1"
6061
},
Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,3 +1,4 @@
11
export * from './definitions'
22
export * from './powersync'
33
export * from './PowerSyncTransactor'
4+
export * from './sqlite-compiler'

‎packages/powersync-db-collection/src/powersync.ts‎

Lines changed: 244 additions & 52 deletions
Original file line numberDiff line numberDiff line change
@@ -1,10 +1,14 @@
11
import { DiffTriggerOperation, sanitizeSQL } from '@powersync/common'
2+
import { Mutex } from 'async-mutex'
3+
import { or } from '@tanstack/db'
4+
import { compileSQLite } from './sqlite-compiler'
25
import { PendingOperationStore } from './PendingOperationStore'
36
import { PowerSyncTransactor } from './PowerSyncTransactor'
47
import { DEFAULT_BATCH_SIZE } from './definitions'
58
import { asPowerSyncRecord, mapOperation } from './helpers'
69
import { convertTableToSchema } from './schema'
710
import { serializeForSQLite } from './serialization'
11+
import type { LoadSubsetOptions, OperationType, SyncConfig } from '@tanstack/db'
812
import type {
913
AnyTableColumnType,
1014
ExtractedTable,
@@ -24,9 +28,8 @@ import type {
2428
PowerSyncCollectionUtils,
2529
} from './definitions'
2630
import type { PendingOperation } from './PendingOperationStore'
27-
import type { SyncConfig } from '@tanstack/db'
2831
import type { StandardSchemaV1 } from '@standard-schema/spec'
29-
import type { Table, TriggerDiffRecord } from '@powersync/common'
32+
import type { LockContext, Table, TriggerDiffRecord } from '@powersync/common'
3033

3134
/**
3235
* Creates PowerSync collection options for use with a standard Collection.
@@ -225,6 +228,7 @@ export function powerSyncCollectionOptions<
225228
table,
226229
schema: inputSchema,
227230
syncBatchSize = DEFAULT_BATCH_SIZE,
231+
syncMode = 'eager',
228232
...restConfig
229233
} = config
230234

@@ -296,11 +300,66 @@ export function powerSyncCollectionOptions<
296300
*/
297301
const sync: SyncConfig<OutputType, string> = {
298302
sync: (params) => {
299-
const { begin, write, commit, markReady } = params
303+
const { begin, write, collection, commit, markReady } = params
300304
const abortController = new AbortController()
301305

302-
// The sync function needs to be synchronous
303-
async function start() {
306+
let disposeTracking: (() => Promise<void>) | null = null
307+
308+
if (syncMode === `eager`) {
309+
return runEagerSync()
310+
} else {
311+
return runOnDemandSync()
312+
}
313+
314+
async function createDiffTrigger(options: {
315+
when: Record<DiffTriggerOperation, string>
316+
writeType: (rowId: string) => OperationType
317+
batchQuery: (
318+
lockContext: LockContext,
319+
batchSize: number,
320+
cursor: number,
321+
) => Promise<Array<TableType>>
322+
onReady: () => void
323+
}) {
324+
const { when, writeType, batchQuery, onReady } = options
325+
326+
return await database.triggers.createDiffTrigger({
327+
source: viewName,
328+
destination: trackedTableName,
329+
when,
330+
hooks: {
331+
beforeCreate: async (context) => {
332+
let currentBatchCount = syncBatchSize
333+
let cursor = 0
334+
while (currentBatchCount == syncBatchSize) {
335+
begin()
336+
337+
const batchItems = await batchQuery(
338+
context,
339+
syncBatchSize,
340+
cursor,
341+
)
342+
currentBatchCount = batchItems.length
343+
cursor += currentBatchCount
344+
for (const row of batchItems) {
345+
write({
346+
type: writeType(row.id),
347+
value: deserializeSyncRow(row),
348+
})
349+
}
350+
commit()
351+
}
352+
onReady()
353+
database.logger.info(
354+
`Sync is ready for ${viewName} into ${trackedTableName}`,
355+
)
356+
},
357+
},
358+
})
359+
}
360+
361+
// The sync function needs to be synchronous.
362+
async function start(afterOnChangeRegistered?: () => Promise<void>) {
304363
database.logger.info(
305364
`Sync is starting for ${viewName} into ${trackedTableName}`,
306365
)
@@ -362,68 +421,200 @@ export function powerSyncCollectionOptions<
362421
},
363422
)
364423

365-
const disposeTracking = await database.triggers.createDiffTrigger({
366-
source: viewName,
367-
destination: trackedTableName,
368-
when: {
369-
[DiffTriggerOperation.INSERT]: `TRUE`,
370-
[DiffTriggerOperation.UPDATE]: `TRUE`,
371-
[DiffTriggerOperation.DELETE]: `TRUE`,
372-
},
373-
hooks: {
374-
beforeCreate: async (context) => {
375-
let currentBatchCount = syncBatchSize
376-
let cursor = 0
377-
while (currentBatchCount == syncBatchSize) {
378-
begin()
379-
const batchItems = await context.getAll<TableType>(
380-
sanitizeSQL`SELECT * FROM ${viewName} LIMIT ? OFFSET ?`,
381-
[syncBatchSize, cursor],
382-
)
383-
currentBatchCount = batchItems.length
384-
cursor += currentBatchCount
385-
for (const row of batchItems) {
386-
write({
387-
type: `insert`,
388-
value: deserializeSyncRow(row),
389-
})
390-
}
391-
commit()
392-
}
393-
markReady()
394-
database.logger.info(
395-
`Sync is ready for ${viewName} into ${trackedTableName}`,
396-
)
397-
},
398-
},
399-
})
424+
await afterOnChangeRegistered?.()
400425

401426
// If the abort controller was aborted while processing the request above
402427
if (abortController.signal.aborted) {
403-
await disposeTracking()
428+
await disposeTracking?.()
404429
} else {
405430
abortController.signal.addEventListener(
406431
`abort`,
407432
() => {
408-
disposeTracking()
433+
disposeTracking?.()
409434
},
410435
{ once: true },
411436
)
412437
}
413438
}
414439

415-
start().catch((error) =>
416-
database.logger.error(
417-
`Could not start syncing process for ${viewName} into ${trackedTableName}`,
418-
error,
419-
),
420-
)
440+
// Eager mode.
441+
// Registers a diff trigger for the entire table.
442+
function runEagerSync() {
443+
start(async () => {
444+
disposeTracking = await createDiffTrigger({
445+
when: {
446+
[DiffTriggerOperation.INSERT]: `TRUE`,
447+
[DiffTriggerOperation.UPDATE]: `TRUE`,
448+
[DiffTriggerOperation.DELETE]: `TRUE`,
449+
},
450+
writeType: (_rowId: string) => `insert`,
451+
batchQuery: (
452+
lockContext: LockContext,
453+
batchSize: number,
454+
cursor: number,
455+
) =>
456+
lockContext.getAll<TableType>(
457+
sanitizeSQL`SELECT * FROM ${viewName} LIMIT ? OFFSET ?`,
458+
[batchSize, cursor],
459+
),
460+
onReady: () => markReady(),
461+
})
462+
}).catch((error) =>
463+
database.logger.error(
464+
`Could not start syncing process for ${viewName} into ${trackedTableName}`,
465+
error,
466+
),
467+
)
468+
469+
return () => {
470+
database.logger.info(
471+
`Sync has been stopped for ${viewName} into ${trackedTableName}`,
472+
)
473+
abortController.abort()
474+
}
475+
}
421476

422-
return () => {
423-
database.logger.info(
424-
`Sync has been stopped for ${viewName} into ${trackedTableName}`,
477+
// On-demand mode.
478+
// Registers a diff trigger for the active WHERE expressions.
479+
function runOnDemandSync() {
480+
start().catch((error) =>
481+
database.logger.error(
482+
`Could not start syncing process for ${viewName} into ${trackedTableName}`,
483+
error,
484+
),
425485
)
426-
abortController.abort()
486+
487+
// Tracks all active WHERE expressions for on-demand sync filtering.
488+
// Each loadSubset call pushes its predicate; unloadSubset removes it.
489+
const activeWhereExpressions: Array<LoadSubsetOptions['where']> = []
490+
const mutex = new Mutex()
491+
492+
const loadSubset = async (options?: LoadSubsetOptions): Promise<void> => {
493+
if (options) {
494+
activeWhereExpressions.push(options.where)
495+
}
496+
497+
if (activeWhereExpressions.length === 0) {
498+
await disposeTracking?.()
499+
return
500+
}
501+
502+
const combinedWhere =
503+
activeWhereExpressions.length === 1
504+
? activeWhereExpressions[0]
505+
: or(
506+
activeWhereExpressions[0]!,
507+
activeWhereExpressions[1]!,
508+
...activeWhereExpressions.slice(2),
509+
)
510+
511+
const compiledNewData = compileSQLite(
512+
{ where: combinedWhere },
513+
{ jsonColumn: 'NEW.data' },
514+
)
515+
516+
const compiledOldData = compileSQLite(
517+
{ where: combinedWhere },
518+
{ jsonColumn: 'OLD.data' },
519+
)
520+
521+
const compiledView = compileSQLite({ where: combinedWhere })
522+
523+
const newDataWhenClause = toInlinedWhereClause(compiledNewData)
524+
const oldDataWhenClause = toInlinedWhereClause(compiledOldData)
525+
const viewWhereClause = toInlinedWhereClause(compiledView)
526+
527+
await disposeTracking?.()
528+
529+
disposeTracking = await createDiffTrigger({
530+
when: {
531+
[DiffTriggerOperation.INSERT]: newDataWhenClause,
532+
[DiffTriggerOperation.UPDATE]: `(${newDataWhenClause}) OR (${oldDataWhenClause})`,
533+
[DiffTriggerOperation.DELETE]: oldDataWhenClause,
534+
},
535+
writeType: (rowId: string) =>
536+
collection.has(rowId) ? `update` : `insert`,
537+
batchQuery: (
538+
lockContext: LockContext,
539+
batchSize: number,
540+
cursor: number,
541+
) =>
542+
lockContext.getAll<TableType>(
543+
`SELECT * FROM ${viewName} WHERE ${viewWhereClause} LIMIT ? OFFSET ?`,
544+
[batchSize, cursor],
545+
),
546+
onReady: () => {},
547+
})
548+
}
549+
550+
const toInlinedWhereClause = (compiled: {
551+
where?: string
552+
params: Array<unknown>
553+
}): string => {
554+
if (!compiled.where) return 'TRUE'
555+
const sqlParts = compiled.where.split('?')
556+
return sanitizeSQL(
557+
sqlParts as unknown as TemplateStringsArray,
558+
...compiled.params,
559+
)
560+
}
561+
562+
const unloadSubset = async (options: LoadSubsetOptions) => {
563+
const idx = activeWhereExpressions.indexOf(options.where)
564+
if (idx !== -1) {
565+
activeWhereExpressions.splice(idx, 1)
566+
}
567+
568+
// Evict rows that were exclusively loaded by the departing predicate.
569+
// These are rows matching the departing WHERE that are no longer covered
570+
// by any remaining active predicate.
571+
const compiledDeparting = compileSQLite({ where: options.where })
572+
const departingWhereSQL = toInlinedWhereClause(compiledDeparting)
573+
574+
let evictionSQL: string
575+
if (activeWhereExpressions.length === 0) {
576+
evictionSQL = `SELECT id FROM ${viewName} WHERE ${departingWhereSQL}`
577+
} else {
578+
const combinedRemaining =
579+
activeWhereExpressions.length === 1
580+
? activeWhereExpressions[0]!
581+
: or(
582+
activeWhereExpressions[0]!,
583+
activeWhereExpressions[1]!,
584+
...activeWhereExpressions.slice(2),
585+
)
586+
const compiledRemaining = compileSQLite({ where: combinedRemaining })
587+
const remainingWhereSQL = toInlinedWhereClause(compiledRemaining)
588+
evictionSQL = `SELECT id FROM ${viewName} WHERE (${departingWhereSQL}) AND NOT (${remainingWhereSQL})`
589+
}
590+
591+
const rowsToEvict = await database.getAll<{ id: string }>(evictionSQL)
592+
if (rowsToEvict.length > 0) {
593+
begin()
594+
for (const { id } of rowsToEvict) {
595+
write({ type: `delete`, key: id })
596+
}
597+
commit()
598+
}
599+
600+
// Recreate the diff trigger for the remaining active WHERE expressions.
601+
await loadSubset()
602+
}
603+
604+
markReady()
605+
606+
return {
607+
cleanup: () => {
608+
database.logger.info(
609+
`Sync has been stopped for ${viewName} into ${trackedTableName}`,
610+
)
611+
abortController.abort()
612+
},
613+
loadSubset: (options: LoadSubsetOptions) =>
614+
mutex.runExclusive(() => loadSubset(options)),
615+
unloadSubset: (options: LoadSubsetOptions) =>
616+
mutex.runExclusive(() => unloadSubset(options)),
617+
}
427618
}
428619
},
429620
// Expose the getSyncMetadata function
@@ -442,6 +633,7 @@ export function powerSyncCollectionOptions<
442633
getKey,
443634
// Syncing should start immediately since we need to monitor the changes for mutations
444635
startSync: true,
636+
syncMode,
445637
sync,
446638
onInsert: async (params) => {
447639
// The transaction here should only ever contain a single insert mutation

0 commit comments

Comments
 (0)