Skip to main content
Wraps an async iterable source (typically from a gRPC server-streaming RPC) Exposes for await, .on(), .filter(), .map(), .take(), .close(), and Symbol.asyncDispose. Only one consumer is allowed per stream instance. Calling .on() or iterating a stream that is already being consumed throws an error. Use .filter() / .map() / .take() to derive new streams before consuming them.

Constructor

Parameters

AsyncIterable<T>
required
() => Promise<void>

Returns

TypedEventStream

Methods

asyncIterator

Returns

AsyncIterator<T>

close()

Signal the stream is done. Interrupts any pending iteration.

Returns

Promise<void>

filter()

Create a filtered sub-stream. When the predicate is a type-guard the returned stream is narrowed to the guard type.

Parameters

(event: T) => event
required

Returns

TypedEventStream
Create a filtered sub-stream. When the predicate is a type-guard the returned stream is narrowed to the guard type.

Parameters

(event: T) => boolean
required

Returns

TypedEventStream

map()

Transform each event into a different shape.

Parameters

(event: T) => U
required

Returns

TypedEventStream

on()

Subscribe to events with a callback. Returns an unsubscribe function that stops the internal consumption loop. The callback may be synchronous or async. If async, back-pressure is applied — the next event is not delivered until the previous callback resolves. Pass onError to handle callback or source errors; otherwise errors are rethrown on the next microtask.

Parameters

(event: T) => void | Promise<void>
required
(error: unknown) => void

Returns

() => void

take()

Yield the first count events, then close the parent stream. count <= 0 returns an empty stream without consuming the parent.

Parameters

number
required

Returns

TypedEventStream