import { BaseObserver, LogLevels, PowerSyncLogger, DBAdapter, Transaction, CrudEntry, CrudBatch, SqliteValue } from '@powersync/common'; import { BucketStorageAdapter, BucketStorageListener, PowerSyncControlCommand, PSInternalTable, rawPowerSyncControl, targetCheckpointRequestId } from './BucketStorageAdapter.js'; import { CrudEntryImpl, CrudEntryJSON } from './CrudEntry.js'; import { MAX_OP_ID } from '../../../constants.js'; /** * @internal */ export class SqliteBucketStorage extends BaseObserver implements BucketStorageAdapter { private updateListener: () => void; private _clientId?: Promise; constructor( private db: DBAdapter, private logger: PowerSyncLogger ) { super(); this.updateListener = db.registerListener({ tablesUpdated: ({ tables }) => { if (tables.includes(PSInternalTable.CRUD)) { this.iterateListeners((l) => l.crudUpdate?.()); } } }); } async dispose() { this.updateListener?.(); } async _getClientId() { const row = await this.db.get<{ client_id: string }>('SELECT powersync_client_id() as client_id'); return row['client_id']; } getClientId() { if (this._clientId == null) { this._clientId = this._getClientId(); } return this._clientId!; } async updateLocalTarget(cb: () => Promise): Promise { const sequenceBefore = await this.db.readTransaction(async (tx): Promise => { const currentCheckpoint = await targetCheckpointRequestId(tx); if (currentCheckpoint != MAX_OP_ID) return; const rs = await tx.getOptional<{ seq: number }>("SELECT seq FROM main.sqlite_sequence WHERE name = 'ps_crud'"); return rs?.seq; }); if (sequenceBefore == null) { // Nothing to update return false; } const opId = await cb(); return this.writeTransaction(async (tx) => { const anyData = await tx.execute('SELECT 1 FROM ps_crud LIMIT 1'); if (anyData.rows?.length) { // if isNotEmpty this.logger.log({ level: LogLevels.debug, message: `New data uploaded since write checkpoint ${opId} - need new write checkpoint` }); return false; } const { seq: seqAfter } = await tx.get<{ seq: number }>( "SELECT seq FROM main.sqlite_sequence WHERE name = 'ps_crud'" ); if (seqAfter != sequenceBefore) { this.logger.log({ level: LogLevels.debug, message: `New data uploaded since write checpoint ${opId} - need new write checkpoint (sequence updated)` }); // New crud data may have been uploaded since we got the checkpoint. Abort. return false; } this.logger.log({ level: LogLevels.debug, message: `Updating target write checkpoint to ${opId}` }); await targetCheckpointRequestId(tx, opId); return true; }); } async nextCrudItem(): Promise { const next = await this.db.getOptional('SELECT * FROM ps_crud ORDER BY id ASC LIMIT 1'); if (!next) { return; } return CrudEntryImpl.fromRow(next); } async hasCrud(): Promise { const anyData = await this.db.getOptional('SELECT 1 FROM ps_crud LIMIT 1'); return !!anyData; } /** * Get a batch of objects to send to the server. * When the objects are successfully sent to the server, call .complete() */ async getCrudBatch(limit: number = 100): Promise { if (!(await this.hasCrud())) { return null; } const crudResult = await this.db.getAll('SELECT * FROM ps_crud ORDER BY id ASC LIMIT ?', [limit]); const all: CrudEntry[] = []; for (const row of crudResult) { all.push(CrudEntryImpl.fromRow(row)); } if (all.length === 0) { return null; } const last = all[all.length - 1]; return { crud: all, haveMore: true, complete: async (writeCheckpoint?: string) => { return this.handleCrudCheckpoint(last.clientId, writeCheckpoint); } }; } handleCrudCheckpoint(lastClientId: number, writeCheckpoint?: string): Promise { return this.writeTransaction(async (tx) => { await tx.execute('DELETE FROM ps_crud WHERE id <= ?', [lastClientId]); if (writeCheckpoint) { const crudResult = await tx.execute('SELECT 1 FROM ps_crud LIMIT 1'); if (crudResult.rows?.length) { await targetCheckpointRequestId(tx, writeCheckpoint); } } else { await targetCheckpointRequestId(tx, MAX_OP_ID); } }); } async writeTransaction(callback: (tx: Transaction) => Promise, options?: { timeoutMs: number }): Promise { return this.db.writeTransaction(callback, options); } async control(op: PowerSyncControlCommand, payload: string | Uint8Array | ArrayBuffer | null): Promise { return await this.writeTransaction(async (tx) => { return (await rawPowerSyncControl(tx, op, payload as SqliteValue))!; }); } async hasMigratedSubkeys(): Promise { const { r } = await this.db.get<{ r: number }>('SELECT EXISTS(SELECT * FROM ps_kv WHERE key = ?) as r', [ SqliteBucketStorage._subkeyMigrationKey ]); return r != 0; } async migrateToFixedSubkeys(): Promise { await this.writeTransaction(async (tx) => { await tx.execute('UPDATE ps_oplog SET key = powersync_remove_duplicate_key_encoding(key);'); await tx.execute('INSERT OR REPLACE INTO ps_kv (key, value) VALUES (?, ?);', [ SqliteBucketStorage._subkeyMigrationKey, '1' ]); }); } static _subkeyMigrationKey = 'powersync_js_migrated_subkeys'; }