Press n or j to go to the next uncovered block, b, p or k for the previous block.
| 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 | 167x 167x 167x 392x 392x 1x 1x 5x 4x 4x 2x 1x 1x 1x 1x 1x 1x 1x 1x | import { HttpErrorResponse } from '@angular/common/http';
import { signal } from '@angular/core';
import { catchError, Observable, Subject, throwError } from 'rxjs';
import { tap } from 'rxjs/operators';
import { Ref } from '../model/ref';
import { printError } from '../util/http';
export type progress = (msg?: string, p?: number) => void;
export type BusEvent = {
event: string,
ref?: Ref,
repost?: Ref,
errors: string[],
};
export class EventBus {
readonly events = new Subject<BusEvent>();
private readonly progressState = signal({
messages: [] as string[],
num: 0,
den: 0,
});
get progressMessages() {
return this.progressState().messages;
}
get progressNum() {
return this.progressState().num;
}
get progressDen() {
return this.progressState().den;
}
private setState(event: string, ref?: Ref, repost?: Ref, errors: string[] = []) {
this.events.next({ event, ref, repost, errors });
console.log('🚌️ Event Bus:', event, event === 'error' ? errors : '', ref);
}
fire(event: string, ref?: Ref, repost?: Ref) {
this.setState(event, ref, repost);
}
fireError(errors: string[], ref?: Ref) {
this.setState('error', ref, undefined, [...errors]);
}
/**
* Download latest revision of ref from the server and then trigger the
* 'refresh' event.
*/
reload(ref?: Ref) {
this.setState('reload', ref);
}
/**
* Notify latest version of ref is not available.
*/
refresh(ref?: Ref) {
this.setState('refresh', ref);
}
/**
* Clear event bus state for sending duplicate events.
*/
reset() {
this.setState('');
}
runAndReload(o: Observable<any>, ref?: Ref) {
return this.runAndReload$(o, ref).subscribe();
}
runAndReload$(o: Observable<any>, ref?: Ref) {
return this.catchError$(o, ref).pipe(tap(() => this.reload(ref)));
}
runAndRefresh(o: Observable<any>, ref?: Ref) {
this.runAndRefresh$(o, ref).subscribe();
}
runAndRefresh$(o: Observable<any>, ref?: Ref) {
return this.catchError$(o, ref).pipe(tap(() => this.refresh(ref)));
}
catchError$(o: Observable<any>, ref?: Ref) {
return o.pipe(
catchError((err: HttpErrorResponse) => {
this.fireError(printError(err), ref);
return throwError(() => err);
})
);
}
isRef(event: BusEvent, r: Ref) {
return event.ref?.url === r.url && event.ref.origin === r.origin;
}
clearProgress(steps = 0) {
if (!steps || !this.progressDen || this.progressNum >= this.progressDen) {
this.progressState.set({ messages: [], num: 0, den: steps });
} else E{
this.progressState.update(progress => ({ ...progress, den: progress.den + steps }));
}
}
msg(msg: string) {
this.progressState.update(progress => ({ ...progress, messages: [...progress.messages, msg] }));
}
steps(steps = 1) {
this.progressState.update(progress => ({ ...progress, den: progress.den + steps }));
}
progress(msg?: string, steps = 1) {
this.progressState.update(progress => ({
...progress,
messages: msg ? [...progress.messages, msg] : progress.messages,
num: progress.num + steps,
}));
}
}
|