• Home
  • Features
  • Pricing
  • Docs
  • Announcements
  • Sign In

bordoley / reactive-js / 13244716502

10 Feb 2025 03:43PM UTC coverage: 98.102% (-0.04%) from 98.141%
13244716502

push

github

bordoley
Change Flowable.create signature, to pass mode as an EventSource

726 of 776 branches covered (93.56%)

Branch coverage included in aggregate %.

6 of 6 new or added lines in 3 files covered. (100.0%)

1 existing line in 1 file now uncovered.

4184 of 4229 relevant lines covered (98.94%)

9630.36 hits per line

Source File
Press 'n' to go to next uncovered line, 'b' for previous

96.61
/src/integrations/node/FlowableStream.ts
1
import { Readable, Transform, Writable } from "stream";
2
import * as Flowable from "../../concurrent/Flowable.js";
1✔
3
import * as Observable from "../../concurrent/Observable.js";
1✔
4
import {
1✔
5
  DeferredObservableWithSideEffectsLike,
6
  DispatcherLike_complete,
7
  FlowableLike,
8
  FlowableLike_flow,
9
  ObservableLike_observe,
10
  PauseableLike_pause,
11
  PauseableLike_resume,
12
} from "../../concurrent.js";
13
import * as EventSource from "../../events/EventSource.js";
1✔
14
import {
1✔
15
  Factory,
16
  Function1,
17
  bindMethod,
18
  ignore,
19
  invoke,
20
  pipe,
21
} from "../../functions.js";
22
import * as Disposable from "../../utils/Disposable.js";
1✔
23
import * as DisposableContainer from "../../utils/DisposableContainer.js";
1✔
24
import {
1✔
25
  DisposableLike,
26
  DisposableLike_dispose,
27
  QueueableLike_backpressureStrategy,
28
  QueueableLike_capacity,
29
  QueueableLike_enqueue,
30
} from "../../utils.js";
31

32
interface FlowableStreamModule {
33
  create(factory: Factory<Readable>): FlowableLike<Uint8Array>;
34

35
  writeTo(
36
    writable: Writable,
37
  ): Function1<
38
    FlowableLike<Uint8Array>,
39
    DeferredObservableWithSideEffectsLike<Uint8Array>
40
  >;
41
}
42

43
type Signature = FlowableStreamModule;
44

45
type NodeStream = Readable | Writable | Transform;
46

47
const disposeStream = (stream: NodeStream) => () => {
12✔
48
  stream.removeAllListeners();
6✔
49
  // Calling destory can result in onError being called
50
  // if we don't catch the error, it crashes the process.
51
  // This kind of sucks, but its the best we can do;
52
  stream.once("error", ignore);
6✔
53
  stream.destroy();
6✔
54
};
55

56
const addToNodeStream =
57
  <TDisposable extends DisposableLike>(
1✔
58
    stream: NodeStream,
59
  ): Function1<TDisposable, TDisposable> =>
60
  disposable => {
3✔
61
    pipe(stream, addDisposable(disposable));
3✔
62
    return disposable;
3✔
63
  };
64

65
const addDisposable =
66
  <TNodeStream extends NodeStream>(
1✔
67
    disposable: DisposableLike,
68
  ): Function1<TNodeStream, TNodeStream> =>
69
  stream => {
9✔
70
    stream.on("error", Disposable.toErrorHandler(disposable));
9✔
71
    stream.once("close", bindMethod(disposable, DisposableLike_dispose));
9✔
72
    pipe(disposable, DisposableContainer.onError(disposeStream(stream)));
9✔
73
    return stream;
9✔
74
  };
75

76
const addToDisposable =
77
  <TNodeStream extends NodeStream>(
1✔
78
    disposable: DisposableLike,
79
  ): Function1<TNodeStream, TNodeStream> =>
80
  stream => {
3✔
81
    pipe(disposable, DisposableContainer.onDisposed(disposeStream(stream)));
3✔
82
    stream.on("error", Disposable.toErrorHandler(disposable));
3✔
83
    return stream;
3✔
84
  };
85

86
export const create: Signature["create"] = factory =>
1✔
87
  Flowable.create(mode =>
3✔
88
    Observable.create<Uint8Array>(observer => {
3✔
89
      const dispatchDisposable = pipe(
3✔
90
        Disposable.create(),
91
        DisposableContainer.onError(Disposable.toErrorHandler(observer)),
92
        DisposableContainer.onComplete(
93
          bindMethod(observer, DispatcherLike_complete),
94
        ),
95
      );
96

97
      const readable = pipe(
3✔
98
        factory(),
99
        addToDisposable(observer),
100
        addDisposable(dispatchDisposable),
101
      );
102

103
      readable.pause();
3✔
104

105
      pipe(
3✔
106
        mode,
107
        EventSource.addEventHandler(isPaused => {
108
          if (isPaused) {
3!
UNCOV
109
            readable.pause();
×
110
          } else {
111
            readable.resume();
3✔
112
          }
113
        }),
114
        addToNodeStream(readable),
115
      );
116

117
      const onData = bindMethod(observer, QueueableLike_enqueue);
3✔
118
      const onEnd = bindMethod(observer, DispatcherLike_complete);
3✔
119

120
      readable.on("data", onData);
3✔
121
      readable.on("end", onEnd);
3✔
122
    }),
123
  );
124

125
export const writeTo: Signature["writeTo"] =
1✔
126
  (
1✔
127
    writable: Writable,
128
  ): Function1<
129
    FlowableLike<Uint8Array>,
130
    DeferredObservableWithSideEffectsLike<Uint8Array>
131
  > =>
132
  flowable =>
3✔
133
    Observable.create<Uint8Array>(observer => {
3✔
134
      pipe(writable, addDisposable(observer));
3✔
135

136
      const flowed = pipe(
3✔
137
        flowable[FlowableLike_flow](observer, {
138
          backpressureStrategy: observer[QueueableLike_backpressureStrategy],
139
          capacity: observer[QueueableLike_capacity],
140
        }),
141
        Disposable.addTo(observer),
142
      );
143

144
      pipe(
3✔
145
        flowed,
146
        Observable.forEach((ev: Uint8Array) => {
147
          // FIXME: when writing to an outgoing node ServerResponse with a UInt8Array
148
          // node throws a type Error regarding expecting a Buffer, though the docs
149
          // say a UInt8Array should be accepted. Need to file a bug.
150
          if (!writable.write(Buffer.from(ev))) {
5✔
151
            flowed[PauseableLike_pause]();
1✔
152
          }
153
        }),
154
        invoke(ObservableLike_observe, observer),
155
      );
156

157
      pipe(
3✔
158
        observer,
159
        DisposableContainer.onComplete(bindMethod(writable, "end")),
160
      );
161

162
      const onDrain = bindMethod(flowed, PauseableLike_resume);
3✔
163
      const onFinish = bindMethod(observer, DisposableLike_dispose);
3✔
164

165
      writable.on("drain", onDrain);
3✔
166
      writable.on("finish", onFinish);
3✔
167

168
      flowed[PauseableLike_resume]();
3✔
169
    });
STATUS · Troubleshooting · Open an Issue · Sales · Support · CAREERS · ENTERPRISE · START FREE TRIAL · SCHEDULE DEMO
ANNOUNCEMENTS · TWITTER · TOS & SLA · Supported CI Services · What's a CI service? · Automated Testing

© 2026 Coveralls, Inc