Angular subscription triggers only once

Viewed 56

I have a material table to show a list of records.

<div class="ml-8 mr-8 mt-8 pb-16">
<div class="mat-elevation-z8">
<mat-table [dataSource]="dataSource$ | async" matSort>
  <ng-container matColumnDef="buildId">
    <mat-header-cell *matHeaderCellDef mat-sort-header> Id </mat-header-cell>
    <mat-cell *matCellDef="let deployment" data-label="buildId"> {{ deployment.buildId }} </mat-cell>
  </ng-container>

  <ng-container matColumnDef="deploymentType">
    <mat-header-cell *matHeaderCellDef mat-sort-header> Type </mat-header-cell>
    <mat-cell *matCellDef="let deployment" data-label="deploymentType"> {{ deployment.deploymentType }} </mat-cell>
  </ng-container>

  <ng-container matColumnDef="status">
    <mat-header-cell *matHeaderCellDef mat-sort-header> Status </mat-header-cell>
    <mat-cell *matCellDef="let deployment" data-label="status"> {{ deployment.status }} </mat-cell>
  </ng-container>

  <ng-container matColumnDef="started">
    <mat-header-cell *matHeaderCellDef mat-sort-header> Started At </mat-header-cell>
    <mat-cell *matCellDef="let deployment" data-label="started">
      {{ deployment.started | date: 'short':locale }}
    </mat-cell>
  </ng-container>

  <ng-container matColumnDef="completed">
    <mat-header-cell *matHeaderCellDef mat-sort-header> Completed At </mat-header-cell>
    <mat-cell *matCellDef="let deployment" data-label="completed">
      {{ deployment.completed | date: 'short':locale }}
    </mat-cell>
  </ng-container>

  <ng-container matColumnDef="queuedCount">
    <mat-header-cell *matHeaderCellDef mat-sort-header> Queued </mat-header-cell>
    <mat-cell *matCellDef="let deployment" data-label="queuedCount"> {{ deployment.queuedCount }} </mat-cell>
  </ng-container>

  <ng-container matColumnDef="completedCount">
    <mat-header-cell *matHeaderCellDef mat-sort-header> Done </mat-header-cell>
    <mat-cell *matCellDef="let deployment" data-label="completedCount"> {{ deployment.completedCount }} </mat-cell>
  </ng-container>

  <ng-container matColumnDef="errorCount">
    <mat-header-cell *matHeaderCellDef mat-sort-header> Error </mat-header-cell>
    <mat-cell *matCellDef="let deployment" data-label="errorCount"> {{ deployment.errorCount }} </mat-cell>
  </ng-container>

  <ng-container matColumnDef="inProgressCount">
    <mat-header-cell *matHeaderCellDef mat-sort-header> In Progress </mat-header-cell>
    <mat-cell *matCellDef="let deployment" data-label="inProgressCount">
      {{ deployment.inProgressCount }}
    </mat-cell>
  </ng-container>

  <ng-container matColumnDef="totalJobCount">
    <mat-header-cell *matHeaderCellDef mat-sort-header> Job Count </mat-header-cell>
    <mat-cell *matCellDef="let deployment" data-label="totalJobCount"> {{ deployment.totalJobCount }} </mat-cell>
  </ng-container>

  <mat-header-row *matHeaderRowDef="displayedColumns"></mat-header-row>
  <mat-row (click)="viewDeploymentDetails(row)" *matRowDef="let row; columns: displayedColumns"> </mat-row>
</mat-table>
<mat-paginator [pageSizeOptions]="[5, 10, 25, 100]" [pageSize]="25"></mat-paginator>

The component subscribes to a service that proxies messages from a websocket connnection

        import { AfterViewInit, ChangeDetectorRef, Component, OnDestroy, OnInit, ViewChild } from '@angular/core';
    import { MatPaginator } from '@angular/material/paginator';
    import { MatSort } from '@angular/material/sort';
    import { MatTableDataSource } from '@angular/material/table';
    import { ActivatedRoute, Router } from '@angular/router';
    import { Deployment } from 'app/shared/models/deployment.model';
    import { DatabaseDeploymentService } from 'app/shared/services/database-deployment.service';
    import { BehaviorSubject, Subject } from 'rxjs';
    import { map } from 'rxjs/operators';

    @Component({
      selector: 'app-list-deployments',
      templateUrl: './list-deployments.component.html',
      styleUrls: ['./list-deployments.component.scss'],
    })
    export class ListDeploymentsComponent implements OnInit, AfterViewInit, OnDestroy {
      @ViewChild(MatPaginator) paginator: MatPaginator;
      @ViewChild(MatSort) sort: MatSort;

      private deploymentsSubject$ = new Subject<MatTableDataSource<Deployment>>();
      private showLoaderSubject$ = new BehaviorSubject<boolean>(true);
      spinnerMode = 'indeterminate';
      displayedColumns: string[] = [
        'buildId',
        'deploymentType',
        'status',
        'started',
        'completed',
        'queuedCount',
        'inProgressCount',
        'completedCount',
        'errorCount',
        'totalJobCount',
      ];
      deployments: Deployment[];
      showLoader$ = this.showLoaderSubject$.asObservable();
      dataSource$ = this.deploymentsSubject$.asObservable();

      constructor(
        private deploymentService: DatabaseDeploymentService,
        private router: Router,
        private route: ActivatedRoute,
        private cdRef: ChangeDetectorRef
      ) {}

      ngOnDestroy(): void {}

      ngOnInit(): void {}

      ngAfterViewInit(): void {
        this.deploymentService.allDeployments$.subscribe((deployments) => {
          this.deployments = deployments;
          this.publishDataSourceUpdate();
          this.showLoaderSubject$.next(false);
        });
        this.deploymentService.deployment$.subscribe((deployment) => {
          this.addOrUpdateDeployments(deployment);
        });
        this.deploymentService.loadDeployments();
        this.cdRef.detectChanges();
      }

      public viewDeploymentDetails(deployment: Deployment) {
        this.router.navigate([`${deployment.buildId}`], { relativeTo: this.route });
      }

      private addOrUpdateDeployments(deployment: Deployment) {
        const found = this.deployments.find((d) => d.id == deployment.id);
        if (found) {
          this.deployments.map((d, i) => {
            if (d.id == deployment.id) {
              this.deployments[i] = deployment;
            }
          });
        } else {
          this.deployments.push(deployment);
        }
        this.publishDataSourceUpdate();
      }

      private publishDataSourceUpdate() {
        const dataSource = new MatTableDataSource<Deployment>();
        dataSource.paginator = this.paginator;
        dataSource.sort = this.sort;
        dataSource.data = this.deployments;
        this.deploymentsSubject$.next(dataSource);
      }
    }

The service implementation:

    import { Injectable } from '@angular/core';
    import { BehaviorSubject, Subject } from 'rxjs';
    import { takeUntil } from 'rxjs/operators';
    import { ApiMessage } from '../models/apiMessage.model';
    import { Deployment } from '../models/deployment.model';
    import { DeploymentDetails } from '../models/deploymentDetails.model';
    import { Job } from '../models/job.model';
    import { RequestMessage } from '../models/requestMessage.model';
    import { WebsocketConnectionService } from './websocket-connection.service';

    @Injectable({
      providedIn: 'root',
    })
    export class DatabaseDeploymentService {
      private messageTypes: string[] = [
        'list-deployments-result',
        'deployment-created',
        'deployment-updated',
        'get-deployment-details-result',
        'job-created',
        'job-updated',
      ];
      private unsubscribe = new Subject<void>();
      private deploymentDetailsSubject$ = new Subject<DeploymentDetails>();
      private deploymentSubject$ = new Subject<Deployment>();
      private allDeploymentsSubject$ = new Subject<Deployment[]>();
      private jobSubject$ = new Subject<Job>();

      
      public deployment$ = this.deploymentSubject$.asObservable();
      public allDeployments$ = this.allDeploymentsSubject$.asObservable();
      public deploymentDetails$ = this.deploymentDetailsSubject$.asObservable();
      public job$ = this.jobSubject$.asObservable();


      constructor(private websocketConnectionService: WebsocketConnectionService) {
        this.websocketConnectionService.messages$
          .pipe(takeUntil(this.unsubscribe))
          .subscribe((apiMessage: ApiMessage) => this.handleMessage(apiMessage));
      }
      ngOnDestroy(): void {
        this.unsubscribe.next();
        this.unsubscribe.complete();
      }

      public loadDeployments() {
        const requestMessage: RequestMessage = {
          action: 'list-deployments',
        };
        this.websocketConnectionService.sendMessage(requestMessage);
      }

      public getDeploymentById(deploymentId: string) {
        const requestMessage: RequestMessage = {
          action: 'get-deployment-details',
          parameter: {
            deploymentId: deploymentId,
          },
        };
        this.websocketConnectionService.sendMessage(requestMessage);
      }

      public getDeploymentByBuild(buildId: string) {
        const requestMessage: RequestMessage = {
          action: 'get-deployment-details',
          parameter: {
            buildId: buildId,
          },
        };
        this.websocketConnectionService.sendMessage(requestMessage);
      }

      private handleMessage(apiMessage: ApiMessage) {
        if (this.messageTypes.includes(apiMessage.type)) {
          switch (apiMessage.type) {
            case 'list-deployments-result':
              this.allDeploymentsSubject$.next(apiMessage.data);
              break;
            case 'deployment-created' || 'deployment-updated':
              this.deploymentSubject$.next(apiMessage.data);
              break;
            case 'get-deployment-details-result':
              this.deploymentDetailsSubject$.next(apiMessage.data);
              break;
            case 'job-created' || 'job-updated':
              this.jobSubject$.next(apiMessage.data);
              break;
            default:
              break;
          }
        }
      }
    }

The issue is my subscription is only being triggered once.

The "deployment$" observable should be triggered multiple times as the updates are occurring but it is not.

I know the messages are coming to the app as I am logging them in the WebSocket.Connection service.

UPDATED.

I added logging statements the switch statement as suggested - the messages are getting published to the subject.

I also added in error handling on the subscription.

The subscription only fires once (I see next in the log). The error handler is never invoked.

For completeness - here is the websocket connection that the deployment service subscribes to:

        import { Injectable } from '@angular/core';
    import { EMPTY, Subject } from 'rxjs';
    import { WebSocketSubject, webSocket } from 'rxjs/webSocket';
    import { catchError, takeUntil, tap } from 'rxjs/operators';
    import { environment } from 'environments/environment';
    import { ApiMessage } from '../models/apiMessage.model';
    import { RequestMessage } from '../models/requestMessage.model';

    const apiUrl = environment.apiURL;

    @Injectable({
      providedIn: 'root',
    })
    export class WebsocketConnectionService {
      private unsubscribe = new Subject<void>();
      private socket$!: WebSocketSubject<any>;
      private messagesSubject$ = new Subject<ApiMessage>();
      public messages$ = this.messagesSubject$.asObservable();

      constructor() {
        this.connect();
      }

      ngOnDestroy(): void {
        this.unsubscribe.next();
        this.unsubscribe.complete();
      }

      public sendMessage(requestMessage: RequestMessage) {
        if (!this.socket$) {
          this.connect();
        }
        const request = JSON.stringify(requestMessage);
        console.log(`Sending ${request}`);
        this.socket$.next(requestMessage);
      }

      private connect() {
        if (!this.socket$ || this.socket$.closed) {
          this.socket$ = this.getNewSocket();
          this.socket$
            .pipe(
              tap({
                error: (error) => console.log(error),
              }),
              catchError((_) => EMPTY),
              takeUntil(this.unsubscribe)
            )
            .subscribe((message) => this.handleMessage(message));
        }
      }

      private getNewSocket() {
        return webSocket({
          url: apiUrl,
          closeObserver: {
            next: () => {
              console.log(`Connection closed`);
            },
          },
        });
      }

      private handleMessage(apiMessage: ApiMessage) {
        if (apiMessage.type) {
          this.messagesSubject$.next(apiMessage);
        }
      }
    }

The updated subscription....

    ngAfterViewInit(): void {
        this.deploymentService.allDeployments$.subscribe((deployments) => {
          this.deployments = deployments;
          this.publishDataSourceUpdate();
          this.showLoaderSubject$.next(false);
        });
        this.deploymentService.deployment$
          .pipe(
            map((deployment) => this.addOrUpdateDeployments(deployment))
          )
          .subscribe(
            (d) => console.log(`next`),  <-- Shows up once
            (e) => console.log(`Error ${e}`, <---never
            () => console.log('complete'))
          );
        this.deploymentService.loadDeployments();
        this.cdRef.detectChanges();
      }
1 Answers

The deployment$ observable emits every time the next method is called on the deploymentSubject. This is done in handleMessage inside your service, which in turn is called whenever this.websocketConnectionService.messages$ emits.

So I would start by checking how many times this one emits. Also, inside handleMessage the deploymentSubject is used only in the deployment-created case. If messages$ emits, I would make sure that it has the correct case value.

Related