All files / app/store bus.ts

61.11% Statements 22/36
80% Branches 12/15
57.14% Functions 16/28
65.51% Lines 19/29

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,
    }));
  }
}