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
9 changes: 9 additions & 0 deletions docs/fsa/fsa-to-fs.md
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,15 @@ const adapter = await FsaNodeSyncAdapterWorker.start('https://<path>/worker.js',
const fs = new FsaNodeFs(dir, adapter);
```

By default, each synchronous operation waits up to 100 milliseconds for the worker's
response. For slower storage, pass a timeout in milliseconds as the third argument:

```js
const adapter = await FsaNodeSyncAdapterWorker.start('https://<path>/worker.js', dir, 1000);
```

The calling thread remains blocked until the worker responds or the timeout is exceeded.

Where `'https://<path>/worker.js'` is a path to a worker file, which could look like this:

```js
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@ export class FsaNodeSyncAdapterWorker implements FsaNodeSyncAdapter {
public static async start(
url: string,
dir: fsa.IFileSystemDirectoryHandle | Promise<fsa.IFileSystemDirectoryHandle>,
timeout?: number,
): Promise<FsaNodeSyncAdapterWorker> {
const worker = new Worker(url);
const future = new Defer<FsaNodeSyncAdapterWorker>();
Expand All @@ -36,7 +37,7 @@ export class FsaNodeSyncAdapterWorker implements FsaNodeSyncAdapter {
switch (code) {
case FsaNodeWorkerMessageCode.Init: {
const [, sab] = msg as FsaNodeWorkerMsgInit;
messenger = new SyncMessenger(sab);
messenger = new SyncMessenger(sab, timeout);
const setRootMessage: FsaNodeWorkerMsgSetRoot = [FsaNodeWorkerMessageCode.SetRoot, id, _dir];
worker.postMessage(setRootMessage);
break;
Expand Down
9 changes: 6 additions & 3 deletions packages/fs-fsa-to-node/src/worker/SyncMessenger.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@ export type AsyncCallback = (request: Uint8Array) => Promise<Uint8Array>;
* @param condition Condition to wait for, when true, the function returns.
* @param ms Maximum time to wait in milliseconds.
*/
const sleepUntil = (condition: () => boolean, ms: number = 100) => {
const sleepUntil = (condition: () => boolean, ms: number) => {
const start = Date.now();
while (!condition()) {
const now = Date.now();
Expand All @@ -31,7 +31,10 @@ export class SyncMessenger {
protected readonly headerSize;
protected readonly dataSize;

public constructor(protected readonly sab: SharedArrayBuffer) {
public constructor(
protected readonly sab: SharedArrayBuffer,
protected readonly timeout: number = 100,
) {
this.int32 = new Int32Array(sab);
this.uint8 = new Uint8Array(sab);
this.headerSize = 4 * 4;
Expand All @@ -46,7 +49,7 @@ export class SyncMessenger {
int32[2] = requestLength;
this.uint8.set(data, headerSize);
Atomics.notify(int32, 0);
sleepUntil(() => int32[1] === 1);
sleepUntil(() => int32[1] === 1, this.timeout);
const responseLength = int32[2];
const response = this.uint8.slice(headerSize, headerSize + responseLength);
return response;
Expand Down
77 changes: 77 additions & 0 deletions packages/fs-fsa-to-node/src/worker/__tests__/SyncMessenger.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,77 @@
import { SyncMessenger } from '../SyncMessenger';
import { FsaNodeSyncAdapterWorker } from '../FsaNodeSyncAdapterWorker';
import { FsaNodeWorkerMessageCode } from '../constants';
import { encoder } from '../../json';
import type { IFileSystemDirectoryHandle } from '@jsonjoy.com/fs-fsa';

describe('synchronous worker timeout', () => {
afterEach(() => {
jest.restoreAllMocks();
});

const respondAfter = (sab: SharedArrayBuffer, delay: number, response: Uint8Array) => {
let elapsed = 0;
jest.spyOn(Date, 'now').mockImplementation(() => {
if (elapsed === delay) {
const header = new Int32Array(sab);
new Uint8Array(sab).set(response, 16);
header[2] = response.length;
header[1] = 1;
}
return elapsed++;
});
};

test('returns a response within the default timeout', () => {
const sab = new SharedArrayBuffer(128);
const response = new Uint8Array([4, 5]);
respondAfter(sab, 50, response);
expect(new SyncMessenger(sab).callSync(new Uint8Array([1, 2, 3]))).toEqual(response);
});

test('keeps the default timeout of 100 milliseconds', () => {
const sab = new SharedArrayBuffer(128);
respondAfter(sab, 200, new Uint8Array([4, 5]));
expect(() => new SyncMessenger(sab).callSync(new Uint8Array([1]))).toThrow('Timeout');
expect(Date.now()).toBe(102);
});

test('allows a response after 100 milliseconds with a longer timeout', () => {
const sab = new SharedArrayBuffer(128);
const response = new Uint8Array([4, 5]);
respondAfter(sab, 200, response);
expect(new SyncMessenger(sab, 500).callSync(new Uint8Array([1, 2, 3]))).toEqual(response);
});

test('still throws when the configured timeout is exceeded', () => {
const sab = new SharedArrayBuffer(128);
respondAfter(sab, 600, new Uint8Array([4, 5]));
expect(() => new SyncMessenger(sab, 500).callSync(new Uint8Array([1]))).toThrow('Timeout');
expect(Date.now()).toBe(502);
});

test('passes the timeout from FsaNodeSyncAdapterWorker.start to synchronous operations', async () => {
const originalWorker = Object.getOwnPropertyDescriptor(globalThis, 'Worker');
const worker = {
onmessage: (_event: { data: unknown[] }) => {},
postMessage: jest.fn(),
};
Object.defineProperty(globalThis, 'Worker', { configurable: true, value: jest.fn(() => worker) });
try {
const sab = new SharedArrayBuffer(128);
const dir = {} as IFileSystemDirectoryHandle;
const started = FsaNodeSyncAdapterWorker.start('worker.js', dir, 500);
await Promise.resolve();
worker.onmessage({ data: [FsaNodeWorkerMessageCode.Init, sab] });
const [, rootId] = worker.postMessage.mock.calls[0][0];
worker.onmessage({ data: [FsaNodeWorkerMessageCode.RootSet, rootId] });
const adapter = await started;
const response = new Uint8Array([4, 5]);
respondAfter(sab, 200, encoder.encode([FsaNodeWorkerMessageCode.Response, response]));
expect(adapter.call('readFile', ['/test.txt'])).toEqual(response);
} finally {
if (originalWorker) Object.defineProperty(globalThis, 'Worker', originalWorker);
else delete (globalThis as any).Worker;
}
});
});
Loading