Skip to main content
Wraps an async iterable source (typically from a gRPC server-streaming RPC) and exposes multiple consumption patterns: for await, .on(), .filter(), .map(), .take(), 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

asyncDispose

Returns

Promise<void>

asyncIterator

Returns

AsyncIterator<T>

close()

Returns

Promise<void>

filter()

Parameters

(event: T) => event
required

Returns

TypedEventStream

Parameters

(event: T) => boolean
required

Returns

TypedEventStream

map()

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.

Parameters

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

Returns

() => void

take()

Parameters

number
required

Returns

TypedEventStream