From c3c392eca34b9c82fc42504e6822fb8c06ddb404 Mon Sep 17 00:00:00 2001 From: Okiki Ojo Date: Sat, 21 Mar 2026 18:11:18 -0400 Subject: [PATCH] feat: enhance observableInputToStream with improved error handling and subscription management Signed-off-by: Okiki Ojo --- helpers/_types.ts | 18 +++++++---- helpers/utils.ts | 78 ++++++++++++++++++++++++++++++++--------------- 2 files changed, 66 insertions(+), 30 deletions(-) diff --git a/helpers/_types.ts b/helpers/_types.ts index e2b3e04..72d9726 100644 --- a/helpers/_types.ts +++ b/helpers/_types.ts @@ -119,22 +119,28 @@ export interface SubscribableLike { } /** - * Observable-like values that can be converted through `Observable.from()`. + * Primary conversion shapes mirrored by helper utilities and interop types. + * + * This alias is intentionally descriptive rather than authoritative. + * `Observable.from()` keeps its own explicit signature so the runtime entry + * points `from()` and `of()` can stay distinct, but the shapes listed here + * mirror the main non-subscribable inputs that `Observable.from()` accepts. */ export type ObservableInputLike = | SpecObservable | AsyncIterable | Iterable - | PromiseLike; + | PromiseLike + | ArrayLike; /** * 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. + * This is intentionally wider than `ObservableInputLike`. The base alias + * covers the primary `Observable.from()`-style conversion set, while interop + * helpers may also need to consume direct subscribables returned by third-party + * libraries such as RxJS. */ export type ObservableInteropInputLike = | ObservableInputLike diff --git a/helpers/utils.ts b/helpers/utils.ts index 6bb6654..2913212 100644 --- a/helpers/utils.ts +++ b/helpers/utils.ts @@ -245,41 +245,69 @@ function hasSubscribe( * 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. + * + * If the foreign subscribable throws during subscription setup, this bridge + * converts that failure into an `ObservableError` value so downstream + * consumers see the same buffered-error behavior as the rest of this library. */ export function observableInputToStream( input: ObservableInteropInputLike, errorContext: string, ): ReadableStream { let subscription: { unsubscribe?(): void } | void; + let settled = false; return new ReadableStream({ start(controller) { const source = hasSubscribe(input) ? input : Observable.from(input); - subscription = source.subscribe({ - next(value) { - try { - controller.enqueue(value); - } catch { - // Downstream cancellation already decided the stream outcome. - } - }, - error(error) { - try { - controller.enqueue(ObservableError.from(error, errorContext)); - controller.close(); - } catch { - // Downstream cancellation already decided the stream outcome. - } - }, - complete() { - try { - controller.close(); - } catch { - // Downstream cancellation already decided the stream outcome. - } - }, - }); + try { + const startedSubscription = source.subscribe({ + next(value) { + try { + controller.enqueue(value); + } catch { + // Downstream cancellation already decided the stream outcome. + } + }, + error(error) { + settled = true; + + try { + controller.enqueue(ObservableError.from(error, errorContext)); + controller.close(); + } catch { + // Downstream cancellation already decided the stream outcome. + } finally { + subscription = undefined; + } + }, + complete() { + settled = true; + + try { + controller.close(); + } catch { + // Downstream cancellation already decided the stream outcome. + } finally { + subscription = undefined; + } + }, + }); + + subscription = settled ? undefined : startedSubscription; + } catch (error) { + settled = true; + + try { + controller.enqueue(ObservableError.from(error, errorContext)); + controller.close(); + } catch { + // Downstream cancellation already decided the stream outcome. + } finally { + subscription = undefined; + } + } }, cancel() { subscription?.unsubscribe?.(); @@ -302,6 +330,8 @@ export function observableInputToStream( * 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. + * That includes synchronous setup failures when the foreign subscribable + * throws from its `subscribe()` implementation. * * @typeParam TIn - Value type accepted by the foreign operator * @typeParam TOut - Value type produced by the foreign operator -- 2.51.2