Skip to content

Commit a1d323b

Browse files
committed
Fix a few minor concurrency issues
1 parent f64dd5d commit a1d323b

14 files changed

Lines changed: 151 additions & 53 deletions
Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
1+
---
2+
'@powersync/common': patch
3+
---
4+
5+
Internal: Fix obtaining crud lock not being abortable.
Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
1+
---
2+
'@powersync/web': patch
3+
---
4+
5+
Fix timeout option having no effect with OPFS WriteAhead VFS.
Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
1+
---
2+
'@powersync/web': patch
3+
---
4+
5+
Log error when a worker fails to load.

packages/common/src/client/AbstractPowerSyncDatabase.ts

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -970,7 +970,7 @@ SELECT * FROM crud_entries;
970970

971971
/**
972972
* Open a read-only transaction.
973-
* Read transactions can run concurrently to a write transaction.
973+
* When multiple connections are available, read transactions can run concurrently to a write transaction.
974974
* Changes from any write transaction are not visible to read transactions started before it.
975975
*
976976
* @param callback - Function to execute within the transaction

packages/common/src/client/sync/stream/AbstractStreamingSyncImplementation.ts

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -379,6 +379,7 @@ export abstract class AbstractStreamingSyncImplementation
379379
private async _uploadAllCrud(signal: AbortSignal): Promise<void> {
380380
return this.obtainLock({
381381
type: LockType.CRUD,
382+
signal,
382383
callback: async () => {
383384
/**
384385
* Keep track of the first item in the CRUD queue for the last `uploadCrud` iteration.

packages/common/src/utils/ControlledExecutor.ts

Lines changed: 9 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -65,14 +65,17 @@ export class ControlledExecutor<T> {
6565

6666
private async execute(param: T) {
6767
this.runningTask = this.task(param);
68-
await this.runningTask;
69-
this.runningTask = undefined;
68+
try {
69+
await this.runningTask;
70+
} finally {
71+
this.runningTask = undefined;
7072

71-
if (this.pendingTaskParam) {
72-
const pendingParam = this.pendingTaskParam;
73-
this.pendingTaskParam = undefined;
73+
if (this.pendingTaskParam) {
74+
const pendingParam = this.pendingTaskParam;
75+
this.pendingTaskParam = undefined;
7476

75-
this.execute(pendingParam);
77+
this.execute(pendingParam);
78+
}
7679
}
7780
}
7881
}

packages/web/src/db/adapters/AsyncWebAdapter.ts

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -126,7 +126,7 @@ function readWritePoolState(writer: DatabaseClient, readers: DatabaseClient[]):
126126
let timeout: any = null;
127127
let release: UnlockFn | undefined;
128128
if (options?.timeoutMs) {
129-
timeout = setTimeout(() => abortController.abort, options.timeoutMs);
129+
timeout = setTimeout(() => abortController.abort('requesting database timed out'), options.timeoutMs);
130130
}
131131

132132
try {

packages/web/src/db/adapters/wa-sqlite/WASQLiteOpenFactory.ts

Lines changed: 17 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -135,16 +135,24 @@ export class WASQLiteOpenFactory implements SQLOpenFactory {
135135
): Promise<DatabaseClient> => {
136136
const workerPort =
137137
typeof optionsDbWorker == 'function'
138-
? resolveWorkerDatabasePortFactory(() =>
139-
optionsDbWorker({
140-
...this.options,
141-
temporaryStorage,
142-
cacheSizeKb,
143-
flags: this.resolvedFlags,
144-
encryptionKey
145-
})
138+
? resolveWorkerDatabasePortFactory(
139+
() =>
140+
optionsDbWorker({
141+
...this.options,
142+
temporaryStorage,
143+
cacheSizeKb,
144+
flags: this.resolvedFlags,
145+
encryptionKey
146+
}),
147+
this.logger
146148
)
147-
: openWorkerDatabasePort(this.options.dbFilename, enableMultiTabs, optionsDbWorker, this.waOptions.vfs);
149+
: openWorkerDatabasePort(
150+
this.options.dbFilename,
151+
enableMultiTabs,
152+
optionsDbWorker,
153+
this.waOptions.vfs,
154+
this.logger
155+
);
148156

149157
const source = Comlink.wrap<OpenWorkerConnection>(workerPort);
150158
const closeSignal = new AbortController();

packages/web/src/db/sync/SharedWebStreamingSyncImplementation.ts

Lines changed: 17 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -5,7 +5,6 @@ import {
55
SyncStatusOptions
66
} from '@powersync/common';
77
import * as Comlink from 'comlink';
8-
import { getNavigatorLocks } from '../../shared/navigator.js';
98
import { AbstractSharedSyncClientProvider } from '../../worker/sync/AbstractSharedSyncClientProvider.js';
109
import { ManualSharedSyncPayload, SharedSyncClientEvent } from '../../worker/sync/SharedSyncImplementation.js';
1110
import { WorkerClient } from '../../worker/sync/WorkerClient.js';
@@ -16,6 +15,7 @@ import {
1615
WebStreamingSyncImplementationOptions
1716
} from './WebStreamingSyncImplementation.js';
1817
import { generateTabCloseSignal } from '../../shared/tab_close_signal.js';
18+
import { logWorkerErrors } from '../../worker/errors.js';
1919

2020
/**
2121
* The shared worker will trigger methods on this side of the message port
@@ -128,22 +128,23 @@ export class SharedWebStreamingSyncImplementation extends WebStreamingSyncImplem
128128
const syncWorker = options.sync?.worker;
129129
if (syncWorker) {
130130
if (typeof syncWorker === 'function') {
131-
this.messagePort = syncWorker(resolvedWorkerOptions).port;
131+
this.messagePort = this.workerPort(syncWorker(resolvedWorkerOptions));
132132
} else {
133-
this.messagePort = new SharedWorker(`${syncWorker}`, {
134-
/* @vite-ignore */
135-
name: `shared-sync-${this.webOptions.identifier}`
136-
}).port;
133+
this.messagePort = this.workerPort(
134+
new SharedWorker(`${syncWorker}`, {
135+
/* @vite-ignore */
136+
name: `shared-sync-${this.webOptions.identifier}`
137+
})
138+
);
137139
}
138140
} else {
139-
this.messagePort = new SharedWorker(
140-
new URL('../../worker/sync/SharedSyncImplementation.worker.js', import.meta.url),
141-
{
141+
this.messagePort = this.workerPort(
142+
new SharedWorker(new URL('../../worker/sync/SharedSyncImplementation.worker.js', import.meta.url), {
142143
/* @vite-ignore */
143144
name: `shared-sync-${this.webOptions.identifier}`,
144145
type: 'module'
145-
}
146-
).port;
146+
})
147+
);
147148
}
148149

149150
/**
@@ -179,6 +180,11 @@ export class SharedWebStreamingSyncImplementation extends WebStreamingSyncImplem
179180
this.isInitialized = this._init();
180181
}
181182

183+
private workerPort(worker: SharedWorker): MessagePort {
184+
logWorkerErrors(worker, this.logger);
185+
return worker.port;
186+
}
187+
182188
protected async _init() {
183189
/**
184190
* The general flow of initialization is:

packages/web/src/worker/db/open-worker-database.ts

Lines changed: 34 additions & 22 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,8 @@
11
import * as Comlink from 'comlink';
22
import { vfsRequiresDedicatedWorkers, WASQLiteVFS } from '../../db/adapters/wa-sqlite/vfs.js';
33
import { OpenWorkerConnection } from '../../db/adapters/wa-sqlite/DatabaseClient.js';
4+
import type { ILogger } from '@powersync/common';
5+
import { logWorkerErrors } from '../errors.js';
46

57
/**
68
* Opens a shared or dedicated worker which exposes opening of database connections
@@ -9,39 +11,45 @@ export function openWorkerDatabasePort(
911
workerIdentifier: string,
1012
multipleTabs = true,
1113
worker: string | URL = '',
12-
vfs?: WASQLiteVFS
14+
vfs?: WASQLiteVFS,
15+
logger?: ILogger
1316
) {
1417
const needsDedicated = vfs && vfsRequiresDedicatedWorkers(vfs);
18+
let resolvedWorker: Worker | SharedWorker;
1519

1620
if (worker) {
17-
return !needsDedicated && multipleTabs
18-
? new SharedWorker(`${worker}`, {
19-
/* @vite-ignore */
20-
name: `shared-DB-worker-${workerIdentifier}`
21-
}).port
22-
: new Worker(`${worker}`, {
23-
/* @vite-ignore */
24-
name: `DB-worker-${workerIdentifier}`
25-
});
21+
resolvedWorker =
22+
!needsDedicated && multipleTabs
23+
? new SharedWorker(`${worker}`, {
24+
/* @vite-ignore */
25+
name: `shared-DB-worker-${workerIdentifier}`
26+
})
27+
: new Worker(`${worker}`, {
28+
/* @vite-ignore */
29+
name: `DB-worker-${workerIdentifier}`
30+
});
2631
} else {
2732
/**
2833
* Webpack V5 can bundle the worker automatically if the full Worker constructor syntax is used
2934
* https://webpack.js.org/guides/web-workers/
3035
* This enables multi tab support by default, but falls back if SharedWorker is not available
3136
* (in the case of Android)
3237
*/
33-
return !needsDedicated && multipleTabs
34-
? new SharedWorker(new URL('./WASQLiteDB.worker.js', import.meta.url), {
35-
/* @vite-ignore */
36-
name: `shared-DB-worker-${workerIdentifier}`,
37-
type: 'module'
38-
}).port
39-
: new Worker(new URL('./WASQLiteDB.worker.js', import.meta.url), {
40-
/* @vite-ignore */
41-
name: `DB-worker-${workerIdentifier}`,
42-
type: 'module'
43-
});
38+
resolvedWorker =
39+
!needsDedicated && multipleTabs
40+
? new SharedWorker(new URL('./WASQLiteDB.worker.js', import.meta.url), {
41+
/* @vite-ignore */
42+
name: `shared-DB-worker-${workerIdentifier}`,
43+
type: 'module'
44+
})
45+
: new Worker(new URL('./WASQLiteDB.worker.js', import.meta.url), {
46+
/* @vite-ignore */
47+
name: `DB-worker-${workerIdentifier}`,
48+
type: 'module'
49+
});
4450
}
51+
52+
return resolveWorkerDatabasePortFactory(() => resolvedWorker, logger);
4553
}
4654

4755
/**
@@ -52,8 +60,12 @@ export function getWorkerDatabaseOpener(workerIdentifier: string, multipleTabs =
5260
return Comlink.wrap<OpenWorkerConnection>(openWorkerDatabasePort(workerIdentifier, multipleTabs, worker));
5361
}
5462

55-
export function resolveWorkerDatabasePortFactory(worker: () => Worker | SharedWorker) {
63+
export function resolveWorkerDatabasePortFactory(worker: () => Worker | SharedWorker, logger?: ILogger) {
5664
const workerInstance = worker();
65+
if (logger) {
66+
logWorkerErrors(workerInstance, logger);
67+
}
68+
5769
return isSharedWorker(workerInstance) ? workerInstance.port : workerInstance;
5870
}
5971

0 commit comments

Comments
 (0)