Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions .changeset/wa-sqlite-opfs-init-fallback.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
'@treecrdt/wa-sqlite': patch
---

Close failed SQLite and OPFS resources, honor memory fallback after initialization failures, and clean up failed worker lifecycles.
6 changes: 5 additions & 1 deletion packages/treecrdt-wa-sqlite/e2e/src/lifecycle.ts
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@ type LifecycleOptions = {
docId: string;
filename: string;
runtime: LifecycleRuntime;
sharedWorkerName?: string;
};

export type LifecycleState = {
Expand All @@ -39,7 +40,10 @@ async function createOpfsLifecycleClient(opts: LifecycleOptions): Promise<Treecr
return createTreecrdtClient({
docId: opts.docId,
storage: { type: 'opfs', filename: opts.filename, fallback: 'throw' },
runtime: { type: opts.runtime },
runtime:
opts.runtime === 'shared-worker' && opts.sharedWorkerName
? { type: 'shared-worker', name: opts.sharedWorkerName }
: { type: opts.runtime },
});
}

Expand Down
77 changes: 60 additions & 17 deletions packages/treecrdt-wa-sqlite/e2e/tests/lifecycle.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,12 @@ import { test, expect, type Page } from '@playwright/test';

type LifecycleHarness = NonNullable<Window['__treecrdtLifecycle']>;
type LifecycleRuntime = 'direct' | 'dedicated-worker' | 'shared-worker';
type LifecycleOptions = {
docId: string;
filename: string;
runtime: LifecycleRuntime;
sharedWorkerName?: string;
};

const scenarios: Array<{
runtime: LifecycleRuntime;
Expand Down Expand Up @@ -39,37 +45,23 @@ async function support(page: Page): Promise<ReturnType<LifecycleHarness['support
});
}

async function drop(
page: Page,
opts: { docId: string; filename: string; runtime: LifecycleRuntime },
) {
async function drop(page: Page, opts: LifecycleOptions) {
await page.evaluate(async (dropOpts) => {
const harness = window.__treecrdtLifecycle;
if (!harness) throw new Error('__treecrdtLifecycle not available');
await harness.drop(dropOpts);
}, opts);
}

async function write(
page: Page,
opts: {
docId: string;
filename: string;
runtime: LifecycleRuntime;
closeBeforeReload?: boolean;
},
) {
async function write(page: Page, opts: LifecycleOptions & { closeBeforeReload?: boolean }) {
return page.evaluate(async (writeOpts) => {
const harness = window.__treecrdtLifecycle;
if (!harness) throw new Error('__treecrdtLifecycle not available');
return await harness.write(writeOpts);
}, opts);
}

async function read(
page: Page,
opts: { docId: string; filename: string; runtime: LifecycleRuntime },
) {
async function read(page: Page, opts: LifecycleOptions) {
return page.evaluate(async (readOpts) => {
const harness = window.__treecrdtLifecycle;
if (!harness) throw new Error('__treecrdtLifecycle not available');
Expand Down Expand Up @@ -143,4 +135,55 @@ test.describe('browser OPFS lifecycle', () => {
});
}
}

test('releases a SharedWorker port after failed OPFS initialization', async ({
page,
}, testInfo) => {
if (testInfo.project.name !== 'chromium-dev') test.skip();
test.setTimeout(120_000);
page.on('console', (msg) => console.log(`[page][${msg.type()}] ${msg.text()}`));

const suffix = `${Date.now().toString(36)}-${Math.random().toString(16).slice(2, 8)}`;
const sharedWorkerName = `lifecycle-recovery-${suffix}`;
const failedOpts: LifecycleOptions = {
docId: `lifecycle-recovery-failed-${suffix}`,
filename: `/${'x'.repeat(512)}.db`,
runtime: 'shared-worker',
sharedWorkerName,
};
const firstOpts: LifecycleOptions = {
docId: `lifecycle-recovery-first-${suffix}`,
filename: `/lifecycle-recovery-first-${suffix}.db`,
runtime: 'shared-worker',
sharedWorkerName,
};
const secondOpts: LifecycleOptions = {
docId: `lifecycle-recovery-second-${suffix}`,
filename: `/lifecycle-recovery-second-${suffix}.db`,
runtime: 'shared-worker',
sharedWorkerName,
};

await waitForHarness(page);
const opfsSupport = await support(page);
if (!opfsSupport.available) test.skip(true, `OPFS unavailable: ${opfsSupport.reason}`);
expect(new TextEncoder().encode(failedOpts.filename).byteLength).toBeGreaterThan(512);

try {
await expect(write(page, failedOpts)).rejects.toThrow(/sqlite3_open_v2|OPFS requested/);

expectReloadedTree(await write(page, { ...firstOpts, closeBeforeReload: true }), {
mode: 'worker',
runtime: 'shared-worker',
});

expectReloadedTree(await write(page, { ...secondOpts, closeBeforeReload: true }), {
mode: 'worker',
runtime: 'shared-worker',
});
} finally {
await drop(page, { ...firstOpts, runtime: 'direct' }).catch(() => {});
await drop(page, { ...secondOpts, runtime: 'direct' }).catch(() => {});
}
});
});
78 changes: 40 additions & 38 deletions packages/treecrdt-wa-sqlite/src/client.ts
Original file line number Diff line number Diff line change
Expand Up @@ -367,24 +367,8 @@ async function createWorkerClient(opts: {
for (const { reject } of pending.values()) reject(err);
pending.clear();
};
worker.addEventListener('message', onMessage);
worker.addEventListener('error', onError);

// init
const initResult = (await call('init', [
opts.baseUrl ?? '/',
opts.filename,
opts.storage,
opts.docId,
])) as { storage?: StorageMode; filename?: string; opfsError?: string } | undefined;
const effectiveStorage: StorageMode = initResult?.storage === 'opfs' ? 'opfs' : 'memory';
const effectiveFilename =
initResult?.filename ??
(effectiveStorage === 'opfs' ? (opts.filename ?? '/treecrdt.db') : ':memory:');
if (effectiveStorage === 'opfs') {
materialized.enableCrossTab({ docId: opts.docId, filename: effectiveFilename });
}
const cleanup = () => {
if (closed) return;
closed = true;
materialized.close();
for (const { reject } of pending.values()) reject(closedError);
Expand All @@ -393,9 +377,23 @@ async function createWorkerClient(opts: {
worker.removeEventListener('message', onMessage);
worker.terminate();
};
worker.addEventListener('message', onMessage);
worker.addEventListener('error', onError);

// init
let initResult: RpcResult<'init'>;
try {
initResult = await call('init', [opts.baseUrl ?? '/', opts.filename, opts.storage, opts.docId]);
} catch (err) {
cleanup();
throw err;
}
const { storage: effectiveStorage, filename: effectiveFilename, opfsError } = initResult;
if (effectiveStorage === 'opfs') {
materialized.enableCrossTab({ docId: opts.docId, filename: effectiveFilename });
}
if (opts.requireOpfs && effectiveStorage !== 'opfs') {
const reason = initResult?.opfsError ? `: ${initResult.opfsError}` : '';
const reason = opfsError ? `: ${opfsError}` : '';
try {
if (!terminalError) await call('close', [] as RpcParams<'close'>);
} catch {
Expand All @@ -410,19 +408,17 @@ async function createWorkerClient(opts: {
if (closed) return;
try {
if (!terminalError) await call('close', [] as RpcParams<'close'>);
cleanup();
} finally {
// noop: cleanup already handles terminal teardown, and repeated close is idempotent
cleanup();
}
};

const dropImpl = async () => {
if (closed) return;
try {
if (!terminalError) await call('drop', [] as RpcParams<'drop'>);
cleanup();
} finally {
// noop: cleanup already handles terminal teardown, and repeated drop is idempotent
cleanup();
}
};

Expand Down Expand Up @@ -512,18 +508,8 @@ async function createSharedWorkerClient(opts: {
for (const { reject } of pending.values()) reject(err);
pending.clear();
};
port.addEventListener('message', onMessage);
port.addEventListener('messageerror', onMessageError);
port.start();

const initResult = (await call('init', [
opts.baseUrl ?? '/',
opts.filename,
opts.storage,
opts.docId,
])) as { storage?: StorageMode; filename?: string; opfsError?: string } | undefined;
const effectiveStorage: StorageMode = initResult?.storage === 'opfs' ? 'opfs' : 'memory';
const cleanup = () => {
if (closed) return;
closed = true;
materialized.close();
for (const { reject } of pending.values()) reject(closedError);
Expand All @@ -532,9 +518,27 @@ async function createSharedWorkerClient(opts: {
port.removeEventListener('messageerror', onMessageError);
port.close();
};
port.addEventListener('message', onMessage);
port.addEventListener('messageerror', onMessageError);
port.start();

let initResult: RpcResult<'init'>;
try {
initResult = await call('init', [opts.baseUrl ?? '/', opts.filename, opts.storage, opts.docId]);
} catch (err) {
try {
if (!terminalError) await call('close', [] as RpcParams<'close'>);
} catch {
// Initialization may not have completed, but the shared worker still needs its port removed.
} finally {
cleanup();
}
throw err;
}
const { storage: effectiveStorage, opfsError } = initResult;

if (opts.requireOpfs && effectiveStorage !== 'opfs') {
const reason = initResult?.opfsError ? `: ${initResult.opfsError}` : '';
const reason = opfsError ? `: ${opfsError}` : '';
try {
if (!terminalError) await call('close', [] as RpcParams<'close'>);
} catch {
Expand All @@ -549,19 +553,17 @@ async function createSharedWorkerClient(opts: {
if (closed) return;
try {
if (!terminalError) await call('close', [] as RpcParams<'close'>);
cleanup();
} finally {
// noop
cleanup();
}
};

const dropImpl = async () => {
if (closed) return;
try {
if (!terminalError) await call('drop', [] as RpcParams<'drop'>);
cleanup();
} finally {
// noop
cleanup();
}
};

Expand Down
5 changes: 3 additions & 2 deletions packages/treecrdt-wa-sqlite/src/node/open.ts
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@ export async function openTreecrdtDbNode(
if (opts.storage === 'opfs' && opts.requireOpfs) {
throw new Error('OPFS is not supported in Node');
}
const loaded = await loadWaSqliteNode(opts.baseUrl);
return openTreecrdtDbFromLoaded({ ...opts, storage: 'memory' }, loaded);
const loadFresh = () => loadWaSqliteNode(opts.baseUrl);
const loaded = await loadFresh();
return openTreecrdtDbFromLoaded({ ...opts, storage: 'memory' }, loaded, loadFresh);
}
Loading
Loading