MCPcopy Create free account
hub / github.com/ReactiveX/rxjs / createPlatformTestObservable

Function createPlatformTestObservable

packages/test/src/test-sources.ts:133–182  ·  view source on GitHub ↗
(
  controller: VirtualTimeController,
  messages: readonly TestMessage<T>[]
)

Source from the content-addressed store, hash-verified

131}
132
133export function createPlatformTestObservable<T>(
134 controller: VirtualTimeController,
135 messages: readonly TestMessage<T>[]
136): TestPlatformObservable<T> {
137 const ObservableConstructor = getObservableConstructor();
138 const subscriptions: MutableSubscriptionLog[] = [];
139 const observable = new ObservableConstructor<T>((subscriber) => {
140 const log = openSubscriptionLog(controller, subscriptions);
141 const activatedAt = controller.now();
142 let closed = false;
143 const close = (): void => {
144 if (!closed) {
145 closed = true;
146 closeSubscriptionLog(controller, log);
147 }
148 };
149 subscriber.addTeardown(close);
150
151 for (const message of messages) {
152 controller.scheduleAt(
153 () => {
154 if (!subscriber.active) {
155 return;
156 }
157 switch (message.notification.kind) {
158 case 'N':
159 subscriber.next(message.notification.value);
160 break;
161 case 'E':
162 subscriber.error(message.notification.error);
163 close();
164 break;
165 case 'C':
166 subscriber.complete();
167 close();
168 break;
169 }
170 },
171 activatedAt + message.frame,
172 { signal: subscriber.signal }
173 );
174 }
175 });
176
177 return Object.assign(observable, {
178 kind: 'observable' as const,
179 messages,
180 subscriptions,
181 });
182}
183
184class IndependentTestSubscriber<T> implements Subscriber<T> {
185 readonly #observer: Partial<Observer<T>>;

Callers 1

observableMethod · 0.85

Calls 9

getObservableConstructorFunction · 0.85
openSubscriptionLogFunction · 0.85
closeFunction · 0.85
scheduleAtMethod · 0.80
nowMethod · 0.65
addTeardownMethod · 0.65
nextMethod · 0.65
errorMethod · 0.65
completeMethod · 0.65

Tested by

no test coverage detected