From 4abf04c12a6acbbfdd5c699fd18316e0973f1e2d Mon Sep 17 00:00:00 2001 From: MathieuG-P <40181755+Zagrios@users.noreply.github.com> Date: Mon, 17 Apr 2023 01:17:19 +0200 Subject: [PATCH] [bugfix-157] create cancellable observable --- .../models/rx/cancellable-observable.class.ts | 42 +++++++++++++++++++ 1 file changed, 42 insertions(+) create mode 100644 src/shared/models/rx/cancellable-observable.class.ts diff --git a/src/shared/models/rx/cancellable-observable.class.ts b/src/shared/models/rx/cancellable-observable.class.ts new file mode 100644 index 00000000..590b4546 --- /dev/null +++ b/src/shared/models/rx/cancellable-observable.class.ts @@ -0,0 +1,42 @@ +import { BehaviorSubject, Observable, Observer, filter, lastValueFrom, take } from "rxjs"; + +export class CancellableObservable extends Observable { + + private readonly _isCancelled$ = new BehaviorSubject(false); + + constructor(subscriber: (obs: CancellableObserver) => void) { + + super((observer: Observer) => { + + const cancellableObserver: CancellableObserver = Object.assign(observer, { + $isCancelled: () => this.$isCancelled(), + onCancel: (callback: () => void) => this.onCancel(callback) + }); + + subscriber(cancellableObserver); + }); + + lastValueFrom(this).finally(() => this._isCancelled$.complete()); + } + + public cancel(): void { + this._isCancelled$.next(true); + } + + public onCancel(callback: () => void): void { + this.$isCancelled().pipe(filter(isCancelled => isCancelled), take(1)).subscribe(callback); + } + + public $isCancelled(): Observable { + return this._isCancelled$.asObservable(); + } + + public isCancelled(): boolean { + return this._isCancelled$.getValue(); + } +} + +export interface CancellableObserver extends Observer { + readonly $isCancelled: () => Observable; + readonly onCancel: (callback: () => void) => void; +} \ No newline at end of file