MCPcopy Create free account
hub / github.com/deanrad/rxfx / wrappedHandler

Function wrappedHandler

service/src/createService.ts:204–245  ·  view source on GitHub ↗
(e: Action<TRequest>)

Source from the content-addressed store, hash-verified

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(

Callers

nothing calls this directly

Calls 7

completeMethod · 0.80
subscribeMethod · 0.80
triggerMethod · 0.80
nextMethod · 0.80
queryMethod · 0.80
unsubscribeMethod · 0.80
handlerFunction · 0.50

Tested by

no test coverage detected