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
Original file line number Diff line number Diff line change
Expand Up @@ -39,7 +39,7 @@ vi.mock('@metamask/streams/browser', async () => {
messageTarget: MockPostMessageTarget;

constructor({ onEnd, messageTarget }: MockStreamOptions) {
super(() => undefined, { readerOnEnd: onEnd, writerOnEnd: onEnd });
super(() => undefined, { onEnd });
MockStream.instances.push(this);
this.messageTarget = messageTarget;
this.messageTarget.onmessage = (event) => {
Expand Down
16 changes: 8 additions & 8 deletions packages/ocap-kernel/src/vats/VatSupervisor.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -40,7 +40,7 @@ const makeVatSupervisor = async ({
platformOptions,
makeAllowedGlobals,
fetchBlob,
writerOnEnd,
onEnd,
}: {
dispatch?: (input: unknown) => void | Promise<void>;
logger?: Logger;
Expand All @@ -49,7 +49,7 @@ const makeVatSupervisor = async ({
platformOptions?: Record<string, unknown>;
makeAllowedGlobals?: (options: { logger: Logger }) => VatEndowments;
fetchBlob?: FetchBlob;
writerOnEnd?: () => void;
onEnd?: () => void;
} = {}): Promise<{
supervisor: VatSupervisor;
stream: TestDuplexStream<JsonRpcMessage, JsonRpcMessage>;
Expand All @@ -59,7 +59,7 @@ const makeVatSupervisor = async ({
JsonRpcMessage
>(dispatch ?? (() => undefined), {
validateInput: isJsonRpcMessage,
writerOnEnd,
onEnd,
});

// Provide a default makePlatform if none is specified
Expand Down Expand Up @@ -160,20 +160,20 @@ describe('VatSupervisor', () => {

it('calls the endowments teardown before closing the stream', async () => {
// The stream is hardened, so we can't vi.spyOn(stream, 'end'). Instead,
// observe the writer's onEnd callback, which fires as part of stream.end().
// observe the stream's onEnd callback, which fires as part of stream.end().
const teardown = vi.fn().mockResolvedValue(undefined);
const writerOnEnd = vi.fn();
const onEnd = vi.fn();
const { supervisor } = await makeVatSupervisor({
makeAllowedGlobals: () => makeVatEndowments({}, teardown),
writerOnEnd,
onEnd,
});

await supervisor.terminate();

expect(teardown).toHaveBeenCalledTimes(1);
expect(writerOnEnd).toHaveBeenCalledTimes(1);
expect(onEnd).toHaveBeenCalledTimes(1);
expect(teardown.mock.invocationCallOrder[0]).toBeLessThan(
writerOnEnd.mock.invocationCallOrder[0] as number,
onEnd.mock.invocationCallOrder[0] as number,
);
});

Expand Down
14 changes: 14 additions & 0 deletions packages/streams/CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,20 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

## [Unreleased]

### Added

- Export `BaseReader` and `BaseWriter` for one-way streams over any transport ([#1138](https://github.com/Consensys-Incorporated/ocap-kernel/pull/1138))

### Changed

- **BREAKING:** `NodePort` requires an `off` method, which `NodeWorkerDuplexStream` uses to remove its listener when the stream ends ([#1138](https://github.com/Consensys-Incorporated/ocap-kernel/pull/1138))
- `split` accepts any number of predicates and narrows the type of each split ([#1138](https://github.com/Consensys-Incorporated/ocap-kernel/pull/1138))
- The `ChromeRuntimeDuplexStream` constructor throws if `localTarget` and `remoteTarget` are the same, as `make()` already did ([#1138](https://github.com/Consensys-Incorporated/ocap-kernel/pull/1138))

### Removed

- **BREAKING:** Remove the `MessagePortReader`, `MessagePortWriter`, `PostMessageReader`, `PostMessageWriter`, `ChromeRuntimeReader`, `ChromeRuntimeWriter`, `NodeWorkerReader`, and `NodeWorkerWriter` exports; use the corresponding duplex streams instead ([#1138](https://github.com/Consensys-Incorporated/ocap-kernel/pull/1138))

### Fixed

- `PostMessageDuplexStream` calls `onEnd` once per stream, and writes after the remote side ends return a done result instead of throwing when `onEnd` closes the transport ([#1137](https://github.com/Consensys-Incorporated/ocap-kernel/pull/1137))
Expand Down
1 change: 0 additions & 1 deletion packages/streams/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -69,7 +69,6 @@
},
"dependencies": {
"@endo/promise-kit": "^1.2.1",
"@endo/stream": "^1.3.1",
"@metamask/kernel-errors": "workspace:^",
"@metamask/kernel-utils": "workspace:^",
"@metamask/superstruct": "^3.4.1",
Expand Down
82 changes: 42 additions & 40 deletions packages/streams/src/BaseDuplexStream.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -167,8 +167,9 @@ describe('BaseDuplexStream', () => {
throw new Error('foo');
});

await expect(stream.synchronize()).rejects.toThrow('foo');
await expect(stream.synchronize()).rejects.toThrow('foo');
const message = 'TestDuplexStream experienced a dispatch failure';
await expect(stream.synchronize()).rejects.toThrow(message);
await expect(stream.synchronize()).rejects.toThrow(message);
});
});

Expand Down Expand Up @@ -343,53 +344,54 @@ describe('BaseDuplexStream', () => {
await sink.return();
});

it('return calls ends both the reader and writer', async () => {
const readerOnEnd = vi.fn();
const writerOnEnd = vi.fn();
const stream = await TestDuplexStream.make(() => undefined, {
readerOnEnd,
writerOnEnd,
});

it.each([
['returning', async (stream: TestDuplexStream) => stream.return()],
[
'throwing',
async (stream: TestDuplexStream) => stream.throw(new Error('foo')),
],
[
'the remote ending',
async (stream: TestDuplexStream) =>
stream.receiveInput(makeStreamDoneSignal()),
],
])('calls onEnd once after %s', async (_, endStream) => {
const onEnd = vi.fn();
const stream = await TestDuplexStream.make(() => undefined, { onEnd });

await endStream(stream);
await stream.return();
expect(readerOnEnd).toHaveBeenCalledOnce();
expect(writerOnEnd).toHaveBeenCalledOnce();
expect(onEnd).toHaveBeenCalledOnce();
expect(await stream.next()).toStrictEqual(makeDoneResult());
});

it('throw calls throw on the writer but return on the reader', async () => {
const readerOnEnd = vi.fn();
const writerOnEnd = vi.fn();
const stream = await TestDuplexStream.make(() => undefined, {
readerOnEnd,
writerOnEnd,
it('ends and calls onEnd if a write fails', async () => {
const onDispatch = vi.fn();
const onEnd = vi.fn();
const stream = await TestDuplexStream.make(onDispatch, { onEnd });
onDispatch.mockImplementation(() => {
throw new Error('foo');
});

await stream.throw(new Error('foo'));
expect(readerOnEnd).toHaveBeenCalledOnce();
expect(writerOnEnd).toHaveBeenCalledOnce();
await expect(stream.write(42)).rejects.toThrow(
'TestDuplexStream experienced a dispatch failure',
);
expect(onEnd).toHaveBeenCalledOnce();
expect(await stream.next()).toStrictEqual(makeDoneResult());
});

it('ending the reader calls reader onEnd function', async () => {
const readerOnEnd = vi.fn();
const stream = await TestDuplexStream.make(() => undefined, {
readerOnEnd,
});
it('dispatches the done signal before calling onEnd', async () => {
const calls: string[] = [];
const stream = await TestDuplexStream.make(
(value) => {
calls.push(stringify(value));
},
{ onEnd: () => calls.push('onEnd') },
);
calls.length = 0;

await stream.receiveInput(makeStreamDoneSignal());
expect(readerOnEnd).toHaveBeenCalledOnce();
});

it('ending the writer calls writer onEnd function', async () => {
const onDispatch = vi.fn(() => {
throw new Error('foo');
});
const writerOnEnd = vi.fn();
const stream = await TestDuplexStream.make(onDispatch, {
writerOnEnd,
});

await expect(stream.write(42)).rejects.toThrow('foo');
expect(writerOnEnd).toHaveBeenCalledOnce();
expect(calls).toStrictEqual([stringify(makeStreamDoneSignal()), 'onEnd']);
});

describe('end', () => {
Expand Down
92 changes: 60 additions & 32 deletions packages/streams/src/BaseDuplexStream.ts
Original file line number Diff line number Diff line change
@@ -1,11 +1,12 @@
import type { PromiseKit } from '@endo/promise-kit';
import { makePromiseKit } from '@endo/promise-kit';
import type { Reader } from '@endo/stream';
import { stringify } from '@metamask/kernel-utils';
import { is, literal, object } from '@metamask/superstruct';
import type { Infer } from '@metamask/superstruct';

import type { BaseReader, BaseWriter, ValidateInput } from './BaseStream.ts';
import { BaseReader, BaseWriter } from './BaseStream.ts';
import type { Dispatch, Listen, OnEnd, ValidateInput } from './BaseStream.ts';
import type { Reader } from './utils.ts';
import { makeDoneResult } from './utils.ts';

export const DuplexStreamSentinel = {
Expand Down Expand Up @@ -46,19 +47,15 @@ export const isDuplexStreamSignal = (
): value is DuplexStreamSignal => isSyn(value) || isAck(value);

/**
* Make a validator for input to a duplex stream. Constructor helper for concrete
* duplex stream implementations.
*
* Validators passed in by consumers must be augmented such that errors aren't
* thrown for {@link DuplexStreamSignal} values.
* Augments a consumer-provided validator so that it accepts
* {@link DuplexStreamSignal} values.
*
* @param validateInput - The validator for the stream's input type.
* @returns A validator for the stream's input type, or `undefined` if no
* validation is desired.
* @returns The augmented validator, or `undefined` if none was provided.
*/
export const makeDuplexStreamInputValidator = <Read>(
const makeDuplexStreamInputValidator = <Read>(
validateInput?: ValidateInput<Read>,
): ((value: unknown) => value is Read) | undefined =>
): ValidateInput<Read> | undefined =>
validateInput &&
((value: unknown): value is Read =>
isDuplexStreamSignal(value) || validateInput(value));
Expand All @@ -77,26 +74,28 @@ const isEnded = (status: SynchronizationStatus): boolean =>
status === SynchronizationStatus.Complete ||
status === SynchronizationStatus.Failed;

/**
* The base of a duplex stream. Essentially a {@link BaseReader} with a `write()` method.
* Backed up by separate {@link BaseReader} and {@link BaseWriter} instances under the hood.
*/
export abstract class BaseDuplexStream<
Read,
ReadStream extends BaseReader<Read>,
Write = Read,
WriteStream extends BaseWriter<Write> = BaseWriter<Write>,
> implements Reader<Read>
{
export type BaseDuplexStreamArgs<Read, Write> = {
name: string;
listen: Listen;
onDispatch: Dispatch<Write>;
validateInput?: ValidateInput<Read> | undefined;
/**
* The underlying reader for the duplex stream.
* Called once when the stream ends, after the final signal has been dispatched.
* For cleanup such as closing the transport.
*/
readonly #reader: ReadStream;
onEnd?: OnEnd | undefined;
};

/**
* The underlying writer for the duplex stream.
*/
readonly #writer: WriteStream;
/**
* The base of a duplex stream over some transport. Essentially a
* {@link BaseReader} with a `write()` method. Backed up by separate
* {@link BaseReader} and {@link BaseWriter} instances under the hood, each of
* which ends the other.
*/
export class BaseDuplexStream<Read, Write = Read> implements Reader<Read> {
readonly #reader: BaseReader<Read>;

readonly #writer: BaseWriter<Write>;

/**
* The promise for the synchronization of the stream with its remote
Expand Down Expand Up @@ -127,10 +126,39 @@ export abstract class BaseDuplexStream<
/**
* Constructs a new {@link BaseDuplexStream}.
*
* @param reader - The underlying reader for the duplex stream.
* @param writer - The underlying writer for the duplex stream.
* @param options - Options bag for configuring the duplex stream.
* @param options.name - The name of the stream, for logging purposes.
* @param options.listen - Subscribes the stream to its transport.
* @param options.onDispatch - Dispatches messages over the transport.
* @param options.validateInput - A function that validates input from the transport.
* @param options.onEnd - A function that is called once when the stream ends.
*/
constructor(reader: ReadStream, writer: WriteStream) {
constructor({
name,
listen,
onDispatch,
validateInput,
onEnd,
}: BaseDuplexStreamArgs<Read, Write>) {
// The writer ends last, so that its final signal is dispatched before onEnd.
const writer: BaseWriter<Write> = new BaseWriter({
name,
onDispatch,
onEnd: async (error) => {
// eslint-disable-next-line @typescript-eslint/no-use-before-define
await reader.return();
await onEnd?.(error);
},
});
const reader = new BaseReader<Read>({
name,
listen,
validateInput: makeDuplexStreamInputValidator(validateInput),
onEnd: async () => {
await writer.return();
},
});

// Set a catch handler to avoid unhandled rejection errors. The promise may
// reject before reads or writes occur, in which case there are no handlers.
this.#syncKit.promise.catch(() => undefined);
Expand Down Expand Up @@ -348,7 +376,7 @@ harden(BaseDuplexStream);
* A duplex stream. Essentially a {@link Reader} with a `write()` method.
*/
export type DuplexStream<Read, Write = Read> = Pick<
BaseDuplexStream<Read, BaseReader<Read>, Write, BaseWriter<Write>>,
BaseDuplexStream<Read, Write>,
'next' | 'write' | 'drain' | 'pipe' | 'return' | 'throw' | 'end'
> & {
[Symbol.asyncIterator]: () => DuplexStream<Read, Write>;
Expand Down
44 changes: 14 additions & 30 deletions packages/streams/src/BaseStream.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -22,13 +22,6 @@ describe('BaseReader', () => {
expect(reader[Symbol.asyncIterator]()).toBe(reader);
});

it('throws if getReceiveInput is called more than once', () => {
const reader = new TestReader();
expect(() => reader.getReceiveInput()).toThrow(
'TestReader received multiple calls to getReceiveInput()',
);
});

it('calls onEnd once when ending', async () => {
const onEnd = vi.fn();
const reader = new TestReader({ onEnd });
Expand Down Expand Up @@ -329,32 +322,23 @@ describe('BaseWriter', () => {
});
});

it('handles repeated failures to dispatch messages', async () => {
const dispatchSpy = vi
.fn()
.mockImplementationOnce(() => {
throw new Error('foo');
})
.mockImplementationOnce(() => {
throw new Error('foo');
});
const writer = new TestWriter({ onDispatch: dispatchSpy });
it('ends the stream if failing to dispatch the error signal', async () => {
const onEnd = vi.fn();
const dispatchSpy = vi.fn(() => {
throw new Error('foo');
});
const writer = new TestWriter({ onDispatch: dispatchSpy, onEnd });

await expect(writer.next(42)).rejects.toThrow(
'TestWriter experienced repeated dispatch failures.',
);
expect(dispatchSpy).toHaveBeenCalledTimes(3);
expect(dispatchSpy).toHaveBeenNthCalledWith(1, 42);
expect(dispatchSpy).toHaveBeenNthCalledWith(2, {
[StreamSentinel.Error]: true,
error: makeErrorMatcher('foo'),
});
expect(dispatchSpy).toHaveBeenNthCalledWith(3, {
[StreamSentinel.Error]: true,
error: makeErrorMatcher(
'TestWriter experienced repeated dispatch failures.',
makeErrorMatcher(
new Error('TestWriter experienced a dispatch failure', {
cause: new Error('foo'),
}),
),
});
);
expect(dispatchSpy).toHaveBeenCalledTimes(2);
expect(onEnd).toHaveBeenCalledOnce();
expect(await writer.next(43)).toStrictEqual(makeDoneResult());
});
});

Expand Down
Loading
Loading