Skip to content
Open
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/shared-worker-terminal-lifecycle.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
'@treecrdt/wa-sqlite': patch
---

Invalidate every shared-worker client when one client drops the shared database, and prune stale ports before resetting the worker session.
10 changes: 10 additions & 0 deletions packages/treecrdt-wa-sqlite/e2e/src/sync.ts
Original file line number Diff line number Diff line change
Expand Up @@ -555,6 +555,14 @@ export async function closeSharedOpfsCrossTabClient(): Promise<void> {
if (client) await client.close();
}

export async function dropSharedOpfsCrossTabClient(): Promise<void> {
sharedOpfsCrossTabUnsubscribe?.();
sharedOpfsCrossTabUnsubscribe = null;
const client = sharedOpfsCrossTabClient;
sharedOpfsCrossTabClient = null;
if (client) await client.drop();
}

export async function runTreecrdtSyncSubscribeE2E(): Promise<{ ok: true }> {
const docId = `e2e-sync-subscribe-${crypto.randomUUID()}`;
const a = await createTreecrdtClient({ storage: memoryStorage, docId });
Expand Down Expand Up @@ -860,6 +868,7 @@ declare global {
__mutateSharedOpfsCrossTabTree?: typeof mutateSharedOpfsCrossTabTree;
__sharedOpfsCrossTabState?: typeof sharedOpfsCrossTabState;
__closeSharedOpfsCrossTabClient?: typeof closeSharedOpfsCrossTabClient;
__dropSharedOpfsCrossTabClient?: typeof dropSharedOpfsCrossTabClient;
}
}

Expand All @@ -874,4 +883,5 @@ if (typeof window !== 'undefined') {
window.__mutateSharedOpfsCrossTabTree = mutateSharedOpfsCrossTabTree;
window.__sharedOpfsCrossTabState = sharedOpfsCrossTabState;
window.__closeSharedOpfsCrossTabClient = closeSharedOpfsCrossTabClient;
window.__dropSharedOpfsCrossTabClient = dropSharedOpfsCrossTabClient;
}
56 changes: 55 additions & 1 deletion packages/treecrdt-wa-sqlite/e2e/tests/cross-tab.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,8 @@ async function waitForCrossTabHarness(page: Page) {
() =>
typeof (window as any).__openSharedOpfsCrossTabClient === 'function' &&
typeof (window as any).__mutateSharedOpfsCrossTabTree === 'function' &&
typeof (window as any).__sharedOpfsCrossTabState === 'function',
typeof (window as any).__sharedOpfsCrossTabState === 'function' &&
typeof (window as any).__dropSharedOpfsCrossTabClient === 'function',
);
}

Expand Down Expand Up @@ -98,6 +99,13 @@ async function closeClient(page: Page) {
});
}

async function dropClient(page: Page) {
await page.evaluate(async () => {
const drop = (window as any).__dropSharedOpfsCrossTabClient;
if (drop) await drop();
});
}

const scenarios: Array<{
name: string;
filePrefix: string;
Expand Down Expand Up @@ -236,3 +244,49 @@ for (const scenario of scenarios) {
}
});
}

test('shared-worker drop invalidates peers and permits a clean replacement', async ({
context,
}, testInfo) => {
if (testInfo.project.name !== 'chromium-dev') test.skip();
test.setTimeout(120_000);

const suffix = `${Date.now().toString(36)}-${Math.random().toString(16).slice(2, 8)}`;
const docId = `e2e-cross-tab-drop-${suffix}`;
const filename = `/e2e-ct-drop-${suffix}.db`;
const pageA = await context.newPage();
const pageB = await context.newPage();
pageA.on('console', (msg) => console.log(`[pageA][${msg.type()}] ${msg.text()}`));
pageB.on('console', (msg) => console.log(`[pageB][${msg.type()}] ${msg.text()}`));

try {
await Promise.all([waitForCrossTabHarness(pageA), waitForCrossTabHarness(pageB)]);
await Promise.all([
openClient(pageA, docId, filename, 'shared-worker'),
openClient(pageB, docId, filename, 'shared-worker'),
]);

const inserted = await mutateTree(pageA, {
replicaLabel: 'cross-tab-drop',
action: 'insert',
nodeInt: 711,
});
await expect
.poll(async () => (await state(pageB)).eventCount, { timeout: 15_000 })
.toBeGreaterThanOrEqual(1);
expect((await state(pageB)).childrenByParent['0'.repeat(32)]).toContain(inserted.node);

await dropClient(pageA);
await expect(state(pageB)).rejects.toThrow('TreeCRDT shared database was dropped');

expect(await openClient(pageB, docId, filename, 'shared-worker')).toEqual({
mode: 'worker',
runtime: 'shared-worker',
storage: 'opfs',
});
expect((await state(pageB)).childrenByParent['0'.repeat(32)]).toEqual([]);
} finally {
await Promise.allSettled([dropClient(pageA), dropClient(pageB)]);
await Promise.allSettled([pageA.close(), pageB.close()]);
}
});
72 changes: 49 additions & 23 deletions packages/treecrdt-wa-sqlite/src/client.ts
Original file line number Diff line number Diff line change
Expand Up @@ -456,8 +456,7 @@ async function createSharedWorkerClient(opts: {
number,
{ resolve: (value: any) => void; reject: (err: Error) => void }
>();
let terminalError: Error | null = null;
let closed = false;
let closedReason: Error | null = null;
let callQueue: Promise<void> = Promise.resolve();

const closedError = new Error(CLIENT_CLOSED_ERROR);
Expand All @@ -468,12 +467,18 @@ async function createSharedWorkerClient(opts: {
);

const callRaw = <M extends RpcMethod>(method: M, params: RpcParams<M>): Promise<RpcResult<M>> => {
if (closed) return Promise.reject(closedError);
if (closedReason) return Promise.reject(closedReason);
const id = nextId++;
if (terminalError) return Promise.reject(terminalError);
return new Promise((resolve, reject) => {
pending.set(id, { resolve, reject });
port.postMessage({ id, method, params } satisfies RpcRequest<M>);
try {
port.postMessage({ id, method, params } satisfies RpcRequest<M>);
} catch (err) {
const postError = new Error(
`shared worker post failed: ${err instanceof Error ? err.message : String(err)}`,
);
cleanup(postError);
}
});
};
const call = <M extends RpcMethod>(method: M, params: RpcParams<M>): Promise<RpcResult<M>> => {
Expand All @@ -483,19 +488,27 @@ async function createSharedWorkerClient(opts: {
};
const materialized = createClientMaterializationDispatcher({
broadcast: (event) => {
if (closed || terminalError) return;
if (closedReason) return;
try {
port.postMessage({ type: 'materialized', event } satisfies RpcPushMessage);
} catch {
// Closing tabs can race a final materialization notification.
} catch (err) {
const postError = new Error(
`shared worker post failed: ${err instanceof Error ? err.message : String(err)}`,
);
cleanup(postError);
}
},
});

const onMessage = (ev: MessageEvent<RpcResponse | RpcPushMessage>) => {
const data = ev.data;
if ('type' in data && data.type === 'materialized') {
materialized.emitIncomingEvent(data.event);
if ('type' in data) {
if (data.type === 'materialized') {
materialized.emitIncomingEvent(data.event);
return;
}
const err = new Error(data.error || 'shared worker terminated');
cleanup(err);
return;
}
const response = data as RpcResponse;
Expand All @@ -507,30 +520,43 @@ async function createSharedWorkerClient(opts: {
};
const onMessageError = () => {
const err = new Error('shared worker message error');
terminalError = err;
for (const { reject } of pending.values()) reject(err);
pending.clear();
try {
port.postMessage({
id: nextId++,
method: 'close',
params: [],
} satisfies RpcRequest<'close'>);
} catch {
// The port may already be unusable; cleanup below is still required.
}
cleanup(err);
};
const cleanup = () => {
if (closed) return;
closed = true;
const onWorkerError = (ev: ErrorEvent) => {
const err = new Error(ev.message || 'shared worker error');
cleanup(err);
};
const cleanup = (reason: Error = closedError) => {
if (closedReason) return;
closedReason = reason;
materialized.close();
for (const { reject } of pending.values()) reject(closedError);
for (const { reject } of pending.values()) reject(reason);
pending.clear();
port.removeEventListener('message', onMessage);
port.removeEventListener('messageerror', onMessageError);
sharedWorker.removeEventListener('error', onWorkerError);
port.close();
};
port.addEventListener('message', onMessage);
port.addEventListener('messageerror', onMessageError);
sharedWorker.addEventListener('error', onWorkerError);
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'>);
if (!closedReason) await call('close', [] as RpcParams<'close'>);
} catch {
// Initialization may not have completed, but the shared worker still needs its port removed.
} finally {
Expand All @@ -543,7 +569,7 @@ async function createSharedWorkerClient(opts: {
if (opts.requireOpfs && effectiveStorage !== 'opfs') {
const reason = opfsError ? `: ${opfsError}` : '';
try {
if (!terminalError) await call('close', [] as RpcParams<'close'>);
if (!closedReason) await call('close', [] as RpcParams<'close'>);
} catch {
// ignore close errors on init failure
} finally {
Expand All @@ -553,18 +579,18 @@ async function createSharedWorkerClient(opts: {
}

const closeImpl = async () => {
if (closed) return;
if (closedReason) return;
try {
if (!terminalError) await call('close', [] as RpcParams<'close'>);
await call('close', [] as RpcParams<'close'>);
} finally {
cleanup();
}
};

const dropImpl = async () => {
if (closed) return;
if (closedReason) return;
try {
if (!terminalError) await call('drop', [] as RpcParams<'drop'>);
await call('drop', [] as RpcParams<'drop'>);
} finally {
cleanup();
}
Expand Down
3 changes: 2 additions & 1 deletion packages/treecrdt-wa-sqlite/src/common-worker.ts
Original file line number Diff line number Diff line change
Expand Up @@ -33,12 +33,13 @@ export class CommonWorkerSession {
}

async closeDbAndReset(): Promise<void> {
if (this.db?.close) await this.db.close();
const db = this.db;
this.db = null;
this.api = null;
this.storedFilename = undefined;
this.storedStorage = 'memory';
this.onAfterReset();
if (db?.close) await db.close();
}

async drop(): Promise<null> {
Expand Down
15 changes: 11 additions & 4 deletions packages/treecrdt-wa-sqlite/src/rpc.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,8 @@
import type { Operation } from '@treecrdt/interface';
import type { MaterializationEvent, MaterializationOutcome } from '@treecrdt/interface/engine';

export const SHARED_WORKER_DROPPED_ERROR = 'TreeCRDT shared database was dropped';

export type RpcStorageMode = 'memory' | 'opfs';

export type RpcSqlParam = number | string | null | Uint8Array;
Expand Down Expand Up @@ -61,10 +63,15 @@ export type RpcResponse<M extends RpcMethod = RpcMethod> =
| { id: number; ok: true; result: RpcResult<M> }
| { id: number; ok: false; error: string };

export type RpcPushMessage = {
type: 'materialized';
event: MaterializationEvent;
};
export type RpcPushMessage =
| {
type: 'materialized';
event: MaterializationEvent;
}
| {
type: 'terminal';
error: string;
};

export function rpcBinaryResult(bytes: Uint8Array | null): Uint8Array | null {
if (bytes === null) return null;
Expand Down
Loading
Loading