| 202 | }; |
| 203 | |
| 204 | const wrappedHandler = (e: Action<TRequest>) => { |
| 205 | const oneResult = handler(e.payload as TRequest); |
| 206 | /* istanbul ignore next */ |
| 207 | const obsResult: Observable<TNext> = |
| 208 | typeof oneResult === 'function' |
| 209 | ? new Observable(oneResult) |
| 210 | : from(oneResult ?? EMPTY); |
| 211 | |
| 212 | const { cancelIdxAtCreation } = e.meta || {}; |
| 213 | delete e.meta; |
| 214 | |
| 215 | return new Observable((observer) => { |
| 216 | // cancelCurrentAndQueued has been called so exit immediately, complete()-ing so others can queue later. |
| 217 | if (cancelCounter.value > cancelIdxAtCreation) { |
| 218 | observer.complete(); |
| 219 | return; |
| 220 | } |
| 221 | |
| 222 | // prettier-ignore |
| 223 | const sub = obsResult |
| 224 | .pipe( |
| 225 | takeUntil(completions.pipe(tap(final => { |
| 226 | if(typeof final !=="undefined") { |
| 227 | bus.trigger(ACs.next(final)) |
| 228 | } |
| 229 | }))), |
| 230 | tap({ |
| 231 | subscribe: () => {bus.trigger(ACs.started(e.payload))}, |
| 232 | complete: () => {bus.trigger(ACs.complete()); observer.complete()}, |
| 233 | unsubscribe: () => {bus.trigger(ACs.canceled()); observer.complete()} |
| 234 | }), |
| 235 | takeUntil(bus.query(ACs.cancel.match)) |
| 236 | ) |
| 237 | .subscribe({ |
| 238 | error(e) { observer.error(e); }, |
| 239 | next(next) { observer.next(next); } |
| 240 | }); |
| 241 | |
| 242 | return () => sub.unsubscribe(); |
| 243 | //@ts-ignore |
| 244 | }).pipe(takeUntil(bus.resets)) as Observable<TNext>; |
| 245 | }; |
| 246 | |
| 247 | /** The main bus listener of this service */ |
| 248 | const mainListener = bus.listen( |