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
14 changes: 11 additions & 3 deletions src/valdi_modules/src/valdi/valdi_core/src/ValdiRuntime.d.ts
Original file line number Diff line number Diff line change
Expand Up @@ -28,14 +28,22 @@ export interface ColorPalette {
[key: string]: string;
}

interface MessageEvent<T = unknown> {
export interface NativeMessageEvent<T = unknown> {
readonly data: T;
readonly ports: readonly NativeMessagePort[];
}

export type OnMessageFunc<T> = (msg: MessageEvent<T>) => void;
export type OnMessageFunc<T> = (msg: NativeMessageEvent<T>) => void;

export interface NativeMessagePort {
onmessage: OnMessageFunc<unknown> | null;
postMessage<T>(data: T, transfer?: readonly NativeMessagePort[]): void;
start(): void;
close(): void;
}

export interface NativeWorker {
postMessage<T>(data: T): void;
postMessage<T>(data: T, transfer?: readonly NativeMessagePort[]): void;
setOnMessage<T>(f: OnMessageFunc<T>): void;
terminate(): void;
}
Expand Down
10 changes: 5 additions & 5 deletions src/valdi_modules/src/valdi/worker/src/Worker.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
import { ValdiRuntime, NativeWorker, OnMessageFunc } from 'valdi_core/src/ValdiRuntime';
import type { NativeMessagePort, NativeWorker, OnMessageFunc, ValdiRuntime } from 'valdi_core/src/ValdiRuntime';

declare const runtime: ValdiRuntime;

Expand All @@ -13,15 +13,15 @@ export class Worker {
this.nativeWorker = runtime.createWorker(url);
}

public set onmessage(func: OnMessageFunc<unknown>) {
public set onmessage(func: (event: MessageEvent<unknown>) => void) {
if (this.nativeWorker) {
this.nativeWorker.setOnMessage(func);
this.nativeWorker.setOnMessage(func as OnMessageFunc<unknown>);
}
}

public postMessage<T>(data: T): void {
public postMessage<T>(data: T, transfer?: readonly MessagePort[]): void {
if (this.nativeWorker) {
this.nativeWorker.postMessage(data);
this.nativeWorker.postMessage(data, transfer as readonly NativeMessagePort[] | undefined);
}
}

Expand Down
168 changes: 168 additions & 0 deletions src/valdi_modules/src/valdi/worker/test/MessageChannel.spec.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,168 @@
import 'jasmine/src/jasmine';

interface NestedTransferredPortData {
readonly ports: readonly MessagePort[];
}

interface TransferredPortData {
readonly nested: NestedTransferredPortData;
}

function receiveNext<T>(port: MessagePort): Promise<MessageEvent<T>> {
return new Promise(resolve => {
port.onmessage = event => resolve(event as MessageEvent<T>);
});
}

function nextTask(): Promise<void> {
return new Promise(resolve => {
// eslint-disable-next-line @snap/valdi/assign-timer-id
setTimeout(resolve, 0);
});
}

describe('MessageChannel', () => {
it('is backed by native MessageChannel and MessagePort classes', () => {
const channel = new MessageChannel();

expect(channel instanceof MessageChannel).toBeTrue();
expect(channel.port1).toBe(channel.port1);
expect(Object.prototype.hasOwnProperty.call(channel, 'port1')).toBeTrue();
expect(Object.prototype.hasOwnProperty.call(channel, 'port2')).toBeTrue();
expect(Object.getOwnPropertyDescriptor(channel, 'port1')?.writable).toBeTrue();
expect(Object.getOwnPropertyDescriptor(channel, 'port2')?.writable).toBeTrue();
expect(Object.prototype.hasOwnProperty.call(channel.port1, 'postMessage')).toBeFalse();
expect(typeof Object.getPrototypeOf(channel.port1).postMessage).toBe('function');
expect(channel.port1.onmessage).toBeNull();

const onmessage = () => {};
channel.port1.onmessage = onmessage;
expect(channel.port1.onmessage).toBe(onmessage);
channel.port1.onmessage = null;
expect(channel.port1.onmessage).toBeNull();
});

it('queues messages until onmessage starts the port and preserves FIFO order', async () => {
const channel = new MessageChannel();
const messages: unknown[] = [];

channel.port1.postMessage('first');
channel.port1.postMessage('second', []);
channel.port1.postMessage('third');

const received = new Promise<void>(resolve => {
channel.port2.onmessage = event => {
messages.push(event.data);
expect(event.ports).toEqual([]);
if (messages.length === 3) {
resolve();
}
};
});

await received;
expect(messages).toEqual(['first', 'second', 'third']);
});

it('delivers ArrayBuffer values', async () => {
const channel = new MessageChannel();
const received = receiveNext<ArrayBuffer>(channel.port2);
const buffer = new Uint8Array([1, 2, 3]).buffer;

channel.port1.postMessage(buffer);

const event = await received;
expect(Array.from(new Uint8Array(event.data))).toEqual([1, 2, 3]);
});

it('discards messages delivered after start when there is no handler', async () => {
const channel = new MessageChannel();
channel.port2.start();
channel.port1.postMessage('discarded');
await nextTask();

const received = receiveNext<string>(channel.port2);
channel.port1.postMessage('delivered');
expect((await received).data).toBe('delivered');
});

it('can transfer a port over another port and queue messages during transfer', async () => {
const control = new MessageChannel();
const payload = new MessageChannel();
const transferredEvent = receiveNext<string>(control.port2);

control.port1.postMessage('transfer', [payload.port1]);
payload.port2.postMessage('queued immediately');

const event = await transferredEvent;
expect(event.data).toBe('transfer');
expect(event.ports.length).toBe(1);
expect((await receiveNext<string>(event.ports[0])).data).toBe('queued immediately');

// Detached handles are inert, including repeated close calls.
payload.port1.postMessage('ignored');
payload.port1.start();
payload.port1.close();
payload.port1.close();
expect(() => control.port1.postMessage('invalid', [payload.port1])).toThrowError(
'MessagePort in transfer list is already detached',
);
});

it('preserves transferred ports embedded in message data', async () => {
const control = new MessageChannel();
const payload = new MessageChannel();
const transferredEvent = receiveNext<TransferredPortData>(control.port2);

control.port1.postMessage({ nested: { ports: [payload.port1, payload.port1] } }, [payload.port1]);

const event = await transferredEvent;
expect(event.ports.length).toBe(1);
expect(event.data.nested.ports[0]).toBe(event.ports[0]);
expect(event.data.nested.ports[1]).toBe(event.ports[0]);

const received = receiveNext<string>(event.data.nested.ports[0]);
payload.port2.postMessage('sent through data port');
expect((await received).data).toBe('sent through data port');
});

it('rejects duplicate transfers atomically', async () => {
const control = new MessageChannel();
const payload = new MessageChannel();

expect(() => control.port1.postMessage('invalid', [payload.port1, payload.port1])).toThrowError(
'Transfer list contains duplicate MessagePort',
);

const received = receiveNext<string>(payload.port2);
payload.port1.postMessage('still attached');
expect((await received).data).toBe('still attached');
});

it('rejects transferring the sending port without detaching it', async () => {
const channel = new MessageChannel();

expect(() => channel.port1.postMessage('invalid', [channel.port1])).toThrowError(
'Transfer list contains source MessagePort',
);

const received = receiveNext<string>(channel.port2);
channel.port1.postMessage('still attached');
expect((await received).data).toBe('still attached');
});

it('closes ports idempotently and discards later messages', async () => {
const channel = new MessageChannel();
let delivered = false;
channel.port2.onmessage = () => {
delivered = true;
};
channel.port2.close();
channel.port2.close();
channel.port2.start();
channel.port1.postMessage('ignored');

await nextTask();
expect(delivered).toBeFalse();
});
});
68 changes: 68 additions & 0 deletions src/valdi_modules/src/valdi/worker/test/WorkerTest.ts
Original file line number Diff line number Diff line change
@@ -1,11 +1,22 @@
import 'jasmine/src/jasmine';
import { Worker, inWorker } from 'worker/src/Worker';

interface WorkerPortMessage {
readonly type: string;
readonly port: MessagePort;
}

function timeout(ms: number): Promise<void> {
// eslint-disable-next-line @snap/valdi/assign-timer-id
return new Promise(resolve => setTimeout(resolve, ms, 'timeout'));
}

function receiveNext<T>(port: MessagePort): Promise<MessageEvent<T>> {
return new Promise(resolve => {
port.onmessage = event => resolve(event as MessageEvent<T>);
});
}

describe('worker', () => {
it('run_with_worker', async () => {
const worker = new Worker('worker/test_workers/TestWorker');
Expand Down Expand Up @@ -47,4 +58,61 @@ describe('worker', () => {
});
expect(await pong).toEqual(echoValue);
}, 100);

it('communicates both ways over a channel transferred to a worker', async () => {
const worker = new Worker('worker/test_workers/MessageChannelWorker');
const channel = new MessageChannel();
const replyEvents: MessageEvent<unknown>[] = [];
const replies = new Promise<void>(resolve => {
channel.port1.onmessage = event => {
replyEvents.push(event);
if (replyEvents.length === 3) {
resolve();
}
};
});

try {
worker.postMessage({ type: 'initialize', port: channel.port2 }, [channel.port2]);
channel.port1.postMessage('first');
channel.port1.postMessage('second');
channel.port1.postMessage('third');

await replies;
expect(replyEvents.map(event => event.data)).toEqual([
['worker reply', 'first', 0, true],
['worker reply', 'second', 0, true],
['worker reply', 'third', 0, true],
]);
expect(replyEvents.map(event => event.ports.length)).toEqual([0, 0, 0]);
} finally {
worker.terminate();
}
});

it('communicates both ways over a channel transferred from a worker', async () => {
const worker = new Worker('worker/test_workers/MessageChannelWorker');
const transferred = new Promise<MessageEvent<unknown>>(resolve => {
worker.onmessage = event => resolve(event);
});

try {
worker.postMessage('transfer from worker');

const event = await transferred;
const data = event.data as WorkerPortMessage;
expect(data.type).toBe('worker channel');
expect(event.ports.length).toBe(1);
expect(data.port).toBe(event.ports[0]);

const port = data.port;
expect((await receiveNext<string>(port)).data).toBe('queued in worker');

const reply = receiveNext<unknown>(port);
port.postMessage('sent to worker');
expect((await reply).data).toEqual(['worker channel reply', 'sent to worker']);
} finally {
worker.terminate();
}
});
});
Original file line number Diff line number Diff line change
@@ -0,0 +1,27 @@
interface WorkerPortMessage {
readonly type: string;
readonly port: MessagePort;
}

onmessage = workerEvent => {
const event = workerEvent as unknown as MessageEvent<unknown>;
if (event.data === 'transfer from worker') {
const channel = new MessageChannel();
const postWorkerMessage = postMessage as (data: unknown, transfer?: readonly MessagePort[]) => void;
channel.port1.onmessage = portEvent => {
channel.port1.postMessage(['worker channel reply', portEvent.data]);
};
postWorkerMessage({ type: 'worker channel', port: channel.port2 }, [channel.port2]);
channel.port1.postMessage('queued in worker');
return;
}

const port = (event.data as WorkerPortMessage).port;
const dataPortMatchesTransferredPort = port === event.ports[0];

// The transferred port becomes the worker's only message receiver.
delete (globalThis as { onmessage?: unknown }).onmessage;
port.onmessage = portEvent => {
port.postMessage(['worker reply', portEvent.data, portEvent.ports.length, dataPortMatchesTransferredPort]);
};
};
1 change: 1 addition & 0 deletions valdi/src/valdi/hermes/HermesJavaScriptCompiler.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -96,6 +96,7 @@ static hermes::hbc::CompileFlags getCompileFlags(hermes::OutputFormatKind output
compileFlags.enableBlockScoping = true;
compileFlags.enableES6Classes = true;
compileFlags.debug = debug;
compileFlags.emitAsyncBreakCheck = true;

return compileFlags;
}
Expand Down
21 changes: 17 additions & 4 deletions valdi/src/valdi/hermes/HermesJavaScriptContext.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@
#include "valdi/runtime/JavaScript/JavaScriptUtils.hpp"

#include "valdi_core/cpp/Text/UTF16Utils.hpp"
#include "valdi_core/cpp/Utils/Defer.hpp"
#include "valdi_core/cpp/Utils/ReferenceInfo.hpp"
#include "valdi_core/cpp/Utils/StaticString.hpp"
#include "valdi_core/cpp/Utils/StringCache.hpp"
Expand Down Expand Up @@ -1393,13 +1394,25 @@ void HermesJavaScriptContext::willEnterVM() {
_enterVMCount++;
}

void HermesJavaScriptContext::requestExecutionTermination() {
IJavaScriptContext::requestExecutionTermination();
_runtime->triggerTimeoutAsyncBreak();
}

void HermesJavaScriptContext::willExitVM(JSExceptionTracker& exceptionTracker) {
SC_ASSERT(_enterVMCount > 0);
--_enterVMCount;
if (_enterVMCount == 0) {
auto result = _runtime->drainJobs();
checkException(result, exceptionTracker);
if (_enterVMCount > 1) {
--_enterVMCount;
return;
}

Valdi::Defer exitVM([this]() { --_enterVMCount; });
if (executionTerminationRequested()) {
return;
}

auto result = _runtime->drainJobs();
checkException(result, exceptionTracker);
}

void HermesJavaScriptContext::startDebugger(bool isWorker) {
Expand Down
1 change: 1 addition & 0 deletions valdi/src/valdi/hermes/HermesJavaScriptContext.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -153,6 +153,7 @@ class HermesJavaScriptContext : public IJavaScriptContext {

void willEnterVM() final;
void willExitVM(JSExceptionTracker& exceptionTracker) final;
void requestExecutionTermination() final;

void startDebugger(bool isWorker) final;
std::optional<IJavaScriptContextDebuggerInfo> getDebuggerInfo() const final;
Expand Down
Loading
Loading