Skip to content

Commit 7343531

Browse files
committed
debug wal
1 parent 865c273 commit 7343531

9 files changed

Lines changed: 103 additions & 61 deletions

File tree

‎packages/shared-internals/src/utils/mutex.ts‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -49,7 +49,7 @@ export class Semaphore<T> {
4949
if (waiter == this.lastWaiter) this.lastWaiter = prev;
5050
}
5151

52-
private requestPermits(amount: number, abort?: AbortSignal): Promise<{ items: T[]; release: UnlockFn }> {
52+
requestPermits(amount: number, abort?: AbortSignal): Promise<{ items: T[]; release: UnlockFn }> {
5353
if (amount <= 0 || amount > this.size) {
5454
throw new Error(`Invalid amount of items requested (${amount}), must be between 1 and ${this.size}`);
5555
}

‎packages/web/package.json‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,10 @@
2323
"react-native": "./dist/index.react_native_web.js",
2424
"default": "./lib/index.js"
2525
},
26+
"./in-memory-wal-experiment": {
27+
"types": "./lib/db/adapters/memory/client.d.ts",
28+
"default": "./lib/db/adapters/memory/client.js"
29+
},
2630
"./bundled_worker": {
2731
"types": "./lib/worker/worker.d.ts",
2832
"default": "./dist/worker/worker.js"

‎packages/web/src/db/adapters/memory/client.ts‎

Lines changed: 41 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
11
import * as Comlink from 'comlink';
22
import { DBAdapter, DBLockOptions, LockContext, RawQueryResult } from '@powersync/common';
3-
import { DatabaseServer, emptyWalState, WalIndexChange, WriteAheadBuffers, WriteAheadState } from './shared.js';
3+
import { applyWalChanges, DatabaseServer, emptyWalState, WalIndexChange, WriteAheadBuffers } from './shared.js';
44
import { Mutex, Semaphore } from '@powersync/shared-internals';
55

66
function createWriteAheadLogBuffers(): WriteAheadBuffers {
@@ -23,6 +23,7 @@ export class InMemoryWriteAheadLogPool extends DBAdapter {
2323
readonly #rawWorkers: PoolWorker[] = [];
2424
readonly #workers: Semaphore<PoolWorker>;
2525
readonly #writeLock = new Mutex();
26+
readonly #walState = emptyWalState();
2627

2728
constructor(options: InMemoryOptions) {
2829
super();
@@ -58,20 +59,50 @@ export class InMemoryWriteAheadLogPool extends DBAdapter {
5859
}, options);
5960
}
6061

62+
#checkpoint() {
63+
const walBuffer = this.#buffers.writeAheadLog;
64+
const databaseBuffer = this.#buffers.database;
65+
const newFileSize = this.#walState.fileSize;
66+
if (databaseBuffer.byteLength < newFileSize) {
67+
databaseBuffer.grow(newFileSize);
68+
}
69+
70+
for (const [pageOffset, overlayEntry] of this.#walState.overlay.entries()) {
71+
const source = new Uint8Array(walBuffer, overlayEntry.logOffset, overlayEntry.size);
72+
new Uint8Array(databaseBuffer, pageOffset).set(source);
73+
}
74+
75+
const cleared: WalIndexChange = { cleared: true, fileSize: newFileSize, walEnd: 0, added: [] };
76+
applyWalChanges(this.#walState, cleared);
77+
for (const worker of this.#rawWorkers) {
78+
worker.addChanges(cleared);
79+
}
80+
}
81+
6182
writeLock<T>(fn: (tx: LockContext) => Promise<T>, options?: DBLockOptions): Promise<T> {
6283
return this.#writeLock.runExclusive(() => {
6384
return this.#withWorker(async (worker) => {
6485
try {
6586
return await fn(worker);
6687
} finally {
6788
const changes = await worker.takeWalChanges();
89+
applyWalChanges(this.#walState, changes);
6890

6991
if (changes.walEnd > 4096 * 128) {
70-
}
71-
72-
for (const otherWorker of this.#rawWorkers) {
73-
if (otherWorker !== worker) {
74-
otherWorker.addChanges(changes);
92+
// Checkpoint. This can't run concurrently to anything else, so acquire remaining workers.
93+
const remainingWorkers = this.#workers.size - 1;
94+
if (remainingWorkers) {
95+
const { release } = await this.#workers.requestPermits(this.#workers.size - 1);
96+
this.#checkpoint();
97+
release();
98+
} else {
99+
this.#checkpoint();
100+
}
101+
} else {
102+
for (const otherWorker of this.#rawWorkers) {
103+
if (otherWorker !== worker) {
104+
otherWorker.addChanges(changes);
105+
}
75106
}
76107
}
77108
}
@@ -113,7 +144,10 @@ class PoolWorker extends LockContext {
113144
constructor(buffers: WriteAheadBuffers) {
114145
super();
115146
this.#buffers = buffers;
116-
this.#worker = new Worker('./worker.js', { type: 'module' });
147+
this.#worker = new Worker(new URL('./worker.js', import.meta.url), { type: 'module' });
148+
this.#worker.onerror = (e) => {
149+
console.error('Worker error', e);
150+
};
117151
this.#server = Comlink.wrap(this.#worker);
118152
}
119153

‎packages/web/src/db/adapters/memory/shared.ts‎

Lines changed: 19 additions & 26 deletions
Original file line numberDiff line numberDiff line change
@@ -14,14 +14,19 @@ export interface WriteAheadState {
1414
/**
1515
* A map of database position offsets to offsets in the WAL logs at which position contents of the page are stored.
1616
*/
17-
overlay: Map<number, number>;
17+
overlay: Map<number, WalOverlayEntry>;
18+
}
19+
20+
export interface WalOverlayEntry {
21+
logOffset: number;
22+
size: number;
1823
}
1924

2025
export interface WalIndexChange {
2126
fileSize: number;
2227
walEnd: number;
2328

24-
added: number[];
29+
added: (number | WalOverlayEntry)[];
2530
cleared: boolean;
2631
}
2732

@@ -37,29 +42,17 @@ export function emptyWalState(): WriteAheadState {
3742
return { fileSize: 0, walEnd: 0, overlay: new Map() };
3843
}
3944

40-
/**
41-
42-
Main tab:
43-
44-
state:
45-
- mutex around writer
46-
- semaphore around readers
45+
export function applyWalChanges(state: WriteAheadState, changes: WalIndexChange) {
46+
state.fileSize = changes.fileSize;
47+
state.walEnd = changes.walEnd;
4748

48-
to start read:
49-
- acquire from semaphore.
50-
- send a wal patch if necessary.
51-
- ... use!
49+
if (changes.cleared) {
50+
state.overlay.clear();
51+
}
5252

53-
to start write:
54-
- acquire from mutex
55-
- ... use, updating write-ahead log offset
56-
- update overlay index if offset has changed
57-
- if len(overlay) > 500
58-
- acquire all readers
59-
- checkpoint, incrementing epoch
60-
61-
Reader:
62-
63-
- obtain overlay from main tab
64-
65-
*/
53+
for (let i = 0; i < changes.added.length; i += 2) {
54+
const dbOffset = changes.added[i] as number;
55+
const walOffset = changes.added[i + 1] as WalOverlayEntry;
56+
state.overlay.set(dbOffset, walOffset);
57+
}
58+
}

‎packages/web/src/db/adapters/memory/vfs.ts‎

Lines changed: 25 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,7 @@
11
// @ts-ignore
22
import { FacadeVFS } from '@journeyapps/wa-sqlite/src/FacadeVFS.js';
33
import * as VFS from '@journeyapps/wa-sqlite/src/VFS.js';
4-
import { emptyWalState, WalIndexChange, WriteAheadBuffers, WriteAheadState } from './shared.js';
4+
import { emptyWalState, WalIndexChange, WalOverlayEntry, WriteAheadBuffers, WriteAheadState } from './shared.js';
55

66
const mainDbSentinel = Symbol();
77

@@ -10,7 +10,7 @@ export class InMemoryWriteAheadLog extends FacadeVFS {
1010
#files = new Map<string, LocalFile | typeof mainDbSentinel>();
1111

1212
#tx: WriteAheadTransaction | undefined = undefined;
13-
#newOverlayPagesToSendToMainTab = new Map<number, number>();
13+
#newOverlayPagesToSendToMainTab = new Map<number, WalOverlayEntry>();
1414

1515
writeAheadState: WriteAheadState = emptyWalState();
1616

@@ -23,7 +23,7 @@ export class InMemoryWriteAheadLog extends FacadeVFS {
2323
}
2424

2525
takeChanges(): WalIndexChange {
26-
const added: number[] = [];
26+
const added: (number | WalOverlayEntry)[] = [];
2727
this.#newOverlayPagesToSendToMainTab.forEach((v, k) => {
2828
added.push(k);
2929
added.push(v);
@@ -40,7 +40,7 @@ export class InMemoryWriteAheadLog extends FacadeVFS {
4040

4141
jAccess(zName: string, _pFlags: number, pResOut: DataView) {
4242
const file = this.#files.get(zName);
43-
pResOut.setInt32(0, file ? 1 : 0, true);
43+
pResOut.setInt32(0, file != null ? 1 : 0, true);
4444
return VFS.SQLITE_OK;
4545
}
4646

@@ -97,7 +97,12 @@ export class InMemoryWriteAheadLog extends FacadeVFS {
9797
const page = overlay.get(offset < 100 ? 0 : offset);
9898

9999
if (page != null) {
100-
const source = new Uint8Array(this.buffers.writeAheadLog, offset < 100 ? page + offset : page, readableBytes);
100+
const pageOffset = page.logOffset;
101+
const source = new Uint8Array(
102+
this.buffers.writeAheadLog,
103+
offset < 100 ? pageOffset + offset : pageOffset,
104+
readableBytes
105+
);
101106
data.set(source);
102107
} else {
103108
// Page is not in WAL overlay, read directly from underlying in-memory buffer.
@@ -114,7 +119,7 @@ export class InMemoryWriteAheadLog extends FacadeVFS {
114119
}
115120
}
116121

117-
if (endOffset < data.byteLength) {
122+
if (readableBytes < data.byteLength) {
118123
data.fill(0, readableBytes); // Fill rest with zeroes.
119124
return VFS.SQLITE_IOERR_SHORT_READ;
120125
}
@@ -179,6 +184,15 @@ export class InMemoryWriteAheadLog extends FacadeVFS {
179184
return VFS.SQLITE_OK;
180185
}
181186

187+
jDelete(name: string) {
188+
if (name === '/database') {
189+
return VFS.SQLITE_IOERR_DELETE;
190+
}
191+
192+
this.#files.delete(name);
193+
return VFS.SQLITE_OK;
194+
}
195+
182196
jDeviceCharacteristics(): number {
183197
return VFS.SQLITE_IOCAP_UNDELETABLE_WHEN_OPEN | VFS.SQLITE_IOCAP_BATCH_ATOMIC;
184198
}
@@ -208,17 +222,17 @@ export class InMemoryWriteAheadLog extends FacadeVFS {
208222
#commit(tx: WriteAheadTransaction) {
209223
this.writeAheadState.walEnd = tx.walEndOffset;
210224
this.writeAheadState.fileSize = tx.fileSize;
211-
tx.changedPages.forEach((walOffset, databaseOffset) => {
212-
this.#newOverlayPagesToSendToMainTab.set(databaseOffset, walOffset);
213-
this.writeAheadState.overlay.set(databaseOffset, walOffset);
225+
tx.changedPages.forEach((walEntry, databaseOffset) => {
226+
this.#newOverlayPagesToSendToMainTab.set(databaseOffset, walEntry);
227+
this.writeAheadState.overlay.set(databaseOffset, walEntry);
214228
});
215229
}
216230
}
217231

218232
class WriteAheadTransaction {
219233
walEndOffset: number;
220234
fileSize: number;
221-
changedPages = new Map<number, number>();
235+
changedPages = new Map<number, WalOverlayEntry>();
222236

223237
constructor(
224238
endOffset: number,
@@ -241,7 +255,7 @@ class WriteAheadTransaction {
241255
new Uint8Array(wal, currentEnd, data.length).set(data);
242256

243257
this.walEndOffset = newEnd;
244-
this.changedPages.set(offset, currentEnd);
258+
this.changedPages.set(offset, { logOffset: currentEnd, size: data.length });
245259
return true;
246260
}
247261

‎packages/web/src/db/adapters/memory/worker.ts‎

Lines changed: 3 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,5 @@
11
import * as Comlink from 'comlink';
2-
import { DatabaseServer, WalIndexChange, WriteAheadBuffers } from './shared.js';
2+
import { applyWalChanges, DatabaseServer, WalIndexChange, WriteAheadBuffers } from './shared.js';
33
import { RawQueryResult } from '@powersync/common';
44
import { InMemoryWriteAheadLog } from './vfs.js';
55
import { RawSqliteConnection } from '../wa-sqlite/RawSqliteConnection.js';
@@ -30,23 +30,11 @@ class MemoryDatabaseServer implements DatabaseServer {
3030

3131
async updateWalState(overlay: WalIndexChange): Promise<void> {
3232
const currentState = this.#vfs.writeAheadState;
33-
currentState.fileSize = overlay.fileSize;
34-
currentState.walEnd = overlay.walEnd;
35-
36-
if (overlay.cleared) {
37-
currentState.overlay.clear();
38-
}
39-
40-
for (let i = 0; i < overlay.added.length; i += 2) {
41-
const dbOffset = overlay.added[i];
42-
const walOffset = overlay.added[i + 1];
43-
currentState.overlay.set(dbOffset, walOffset);
44-
}
33+
applyWalChanges(currentState, overlay);
4534
}
4635

4736
async executeRaw(query: string, params?: any[] | undefined): Promise<RawQueryResult> {
48-
const results = await this.#connection.executeRaw(query, params);
49-
return results[0];
37+
return await this.#connection.execute(query, params);
5038
}
5139

5240
async takeWalChanges(): Promise<WalIndexChange> {

‎packages/web/src/db/adapters/wa-sqlite/RawSqliteConnection.ts‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -61,7 +61,7 @@ export class RawSqliteConnection {
6161
}
6262

6363
async initWithModule(module: any, vfs: any) {
64-
const api = await this.openSQLiteAPI(module, vfs);
64+
const api = (this._sqliteAPI = await this.openSQLiteAPI(module, vfs));
6565
this.db = await api.open_v2(
6666
this.options.filename,
6767
this.options.readonly ? 1 /* SQLITE_OPEN_READONLY */ : 6 /* SQLITE_OPEN_READWRITE | SQLITE_OPEN_CREATE */

‎packages/web/tests/main.test.ts‎

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,5 @@
11
import { PowerSyncDatabase, WASQLiteOpenFactory, WASQLiteVFS } from '@powersync/web';
2+
import { InMemoryWriteAheadLogPool } from '@powersync/web/in-memory-wal-experiment';
23
import { v4 as uuid } from 'uuid';
34
import { describe, expect, it } from 'vitest';
45
import { TEST_SCHEMA, TestDatabase } from './utils/test-schema.js';
@@ -94,6 +95,13 @@ describe('Basic - with in-memory', () => {
9495
})
9596
)
9697
);
98+
99+
describe(
100+
'wal',
101+
describeBasicTests(() =>
102+
generateTestDb({ schema: TEST_SCHEMA, opened: new InMemoryWriteAheadLogPool({ numWorkers: 1 }) })
103+
)
104+
);
97105
});
98106

99107
function describeBasicTests(generateDB: () => PowerSyncDatabase) {

‎packages/web/vitest.config.ts‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,7 @@ const config: ViteUserConfigExport = {
1717
* first. This is required due to the format of Webworker URIs
1818
* they link to `.js` files.
1919
*/
20+
'@powersync/web/in-memory-wal-experiment': path.resolve(__dirname, './lib/db/adapters/memory/client.js'),
2021
'@powersync/web': path.resolve(__dirname, './lib'),
2122
// Mock WebRemote to throw 401 errors for all HTTP requests in tests
2223
'../../db/sync/WebRemote.js': path.resolve(__dirname, './tests/mocks/MockWebRemote.ts')

0 commit comments

Comments
 (0)