diff --git a/helpers/_types.ts b/helpers/_types.ts index 5c114bf..e2b3e04 100644 --- a/helpers/_types.ts +++ b/helpers/_types.ts @@ -106,6 +106,18 @@ export interface StreamPair { writable: WritableStream; } +/** + * Minimal subscribable shape used by foreign Observable implementations. + */ +export interface SubscribableLike { + /** Subscribe to values from the foreign Observable-like object. */ + subscribe(observer: { + next?(value: T): void; + error?(error: unknown): void; + complete?(): void; + }): { unsubscribe?(): void } | void; +} + /** * Observable-like values that can be converted through `Observable.from()`. */ @@ -115,13 +127,51 @@ export type ObservableInputLike = | Iterable | PromiseLike; +/** + * Wider Observable-like values accepted by interop helpers that adapt foreign + * operator ecosystems. + * + * This is intentionally wider than `ObservableInputLike`. `Observable.from()` + * stays aligned with the library's spec-facing conversion contract, while + * interop helpers may also need to consume direct subscribables returned by + * third-party libraries such as RxJS. + */ +export type ObservableInteropInputLike = + | ObservableInputLike + | SubscribableLike; + /** * Foreign operator shape used by libraries that transform one Observable-like * source into another Observable-like result. */ -export type ObservableOperatorInterop = ( - source: SpecObservable, -) => ObservableInputLike; +export type ObservableOperatorInterop< + TIn, + TOut, + TSource = Observable, +> = ( + source: TSource, +) => ObservableInteropInputLike; + +/** + * Configuration for adapting foreign Observable-style operators. + * + * Some libraries, such as RxJS, require their own Observable class on the + * input side even when the returned operator shape is still `(source) => + * output`. `sourceAdapter` lets callers convert this library's Observable into + * that foreign source type before the operator runs. + */ +export interface ObservableOperatorInteropOptions> { + /** + * Converts this library's Observable into the source type expected by the + * foreign operator. + */ + sourceAdapter?: (source: Observable) => TSource; + + /** + * Context label used when the foreign output reports an error. + */ + errorContext?: string; +} // ======================================== // 2. CREATEOPERATOR INTERFACES diff --git a/helpers/utils.ts b/helpers/utils.ts index 14497a4..2514df9 100644 --- a/helpers/utils.ts +++ b/helpers/utils.ts @@ -1,6 +1,7 @@ import type { - ObservableInputLike, + ObservableInteropInputLike, ObservableOperatorInterop, + ObservableOperatorInteropOptions, TransformFunctionOptions, StreamPair, TransformStreamOptions, @@ -217,22 +218,43 @@ export function streamAsObservable(stream: ReadableStream): Observable }); } +/** + * Returns true when a value exposes a direct `subscribe()` method. + * + * Many Observable libraries return subscribable objects that are usable without + * first going through this library's `Observable.from()`. Detecting that shape + * lets interop stay direct for outputs such as RxJS Observables. + */ +function hasSubscribe( + input: ObservableInteropInputLike | unknown, +): input is { + subscribe(observer: { + next?(value: T): void; + error?(error: unknown): void; + complete?(): void; + }): { unsubscribe?(): void } | void; +} { + return typeof input === "object" && input !== null && + typeof (input as { subscribe?: unknown }).subscribe === "function"; +} + /** * Subscribes to an Observable-like output and exposes it as a ReadableStream. * - * Foreign operators are allowed to return any shape that `Observable.from()` - * understands. Converting the result with a direct subscription avoids the - * extra async-generator and stream layers that `pull(...)+toStream()` would add. + * Foreign operators are allowed to return either a normal `Observable.from()` + * input or a direct subscribable from another Observable implementation. + * Converting the result with a direct subscription avoids the extra + * async-generator and stream layers that `pull(...)+toStream()` would add. */ export function observableInputToStream( - input: ObservableInputLike, + input: ObservableInteropInputLike, errorContext: string, ): ReadableStream { return new ReadableStream({ start(controller) { - const observable = Observable.from(input); + const source = hasSubscribe(input) ? input : Observable.from(input); - return observable.subscribe({ + return source.subscribe({ next(value) { try { controller.enqueue(value); @@ -271,9 +293,9 @@ export function observableInputToStream( * 2. calling the foreign operator * 3. converting the resulting Observable-like output back into a stream * - * The result keeps this library's buffered error behavior by reading the - * foreign output through `pull(..., { throwError: false })` before converting it - * back to a `ReadableStream`. + * The result keeps this library's buffered error behavior by converting the + * foreign output into a `ReadableStream` that enqueues wrapped + * `ObservableError` values instead of failing the readable side outright. * * @typeParam TIn - Value type accepted by the foreign operator * @typeParam TOut - Value type produced by the foreign operator @@ -284,19 +306,24 @@ export function observableInputToStream( * ```ts * const foreignTakeOne = fromObservableOperator((source) => * rxTake(1)(source) - * ); + * , { sourceAdapter: (source) => rxFrom(source) }); * ``` */ -export function fromObservableOperator( - operator: ObservableOperatorInterop, +export function fromObservableOperator>( + operator: ObservableOperatorInterop, + options?: ObservableOperatorInteropOptions, ): Operator { return (source) => { const observableSource = streamAsObservable(source); - const output = operator(observableSource) as ObservableInputLike; + const foreignSource = options?.sourceAdapter + ? options.sourceAdapter(observableSource) + : observableSource as unknown as TSource; + + const output = operator(foreignSource) as ObservableInteropInputLike; return observableInputToStream( output, - "operator:fromObservableOperator:output", + options?.errorContext ?? "operator:fromObservableOperator:output", ); }; }