import { EventStreamMarshaller as UniversalEventStreamMarshaller } from "@smithy/eventstream-serde-universal"; import { Decoder, Encoder, EventStreamMarshaller as IEventStreamMarshaller, Message } from "@smithy/types"; import { Readable } from "stream"; import { readabletoIterable } from "./utils"; /** * @internal */ export interface EventStreamMarshaller extends IEventStreamMarshaller {} /** * @internal */ export interface EventStreamMarshallerOptions { utf8Encoder: Encoder; utf8Decoder: Decoder; } /** * @internal */ export class EventStreamMarshaller { private readonly universalMarshaller: UniversalEventStreamMarshaller; constructor({ utf8Encoder, utf8Decoder }: EventStreamMarshallerOptions) { this.universalMarshaller = new UniversalEventStreamMarshaller({ utf8Decoder, utf8Encoder, }); } deserialize(body: Readable, deserializer: (input: Record) => Promise): AsyncIterable { //should use stream[Symbol.asyncIterable] when the api is stable //reference: https://nodejs.org/docs/latest-v11.x/api/stream.html#stream_readable_symbol_asynciterator const bodyIterable: AsyncIterable = typeof body[Symbol.asyncIterator] === "function" ? body : readabletoIterable(body); return this.universalMarshaller.deserialize(bodyIterable, deserializer); } serialize(input: AsyncIterable, serializer: (event: T) => Message): Readable { return Readable.from(this.universalMarshaller.serialize(input, serializer)); } }