( controller: VirtualTimeController, messages: readonly TestMessage<T>[] )
| 131 | } |
| 132 | |
| 133 | export 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 | |
| 184 | class IndependentTestSubscriber<T> implements Subscriber<T> { |
| 185 | readonly #observer: Partial<Observer<T>>; |
no test coverage detected