All files MessageProcessor.ts

100% Statements 17/17
57.14% Branches 4/7
100% Functions 3/3
100% Lines 17/17

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 811x   1x                                                               1x       1x             2x             3x         3x 3x       1x       1x 1x   1x 1x 1x   1x 1x 1x          
import { Subject, delay } from "rxjs";
import { IEventResult } from "./Event";
import {
  ProcessErrorArgs,
  ServiceBusError,
  ServiceBusErrorCode,
  ServiceBusReceivedMessage,
  isServiceBusError,
} from "@azure/service-bus";
import { ILogger } from "./ServiceBusPubSub";
import { error } from "console";
 
export interface IMessageProcessor {
  /**
   *
   * @param subStream Rxjs channel based stream, each channel/client would have it's own stream pipeline.
   * @param message ServiceBusReceivedMessage
   */
  process(
    subStream: Subject<IEventResult>,
    message: ServiceBusReceivedMessage
  ): Promise<void>;
 
  /**
   *
   * @param args ServiceBus ProcessErrorArgs
   *
   */
  onError(
    args: ProcessErrorArgs,
    handleError: (error: ServiceBusError) => Promise<void>
  ): Promise<void>;
}
 
export class MessageProcessor implements IMessageProcessor {
  private logger: ILogger;
 
  constructor(logger: ILogger) {
    this.logger = logger;
  }
 
  async process(
    subject: Subject<IEventResult>,
    message: ServiceBusReceivedMessage
  ): Promise<void> {
    await subject.next({ ...message });
  }
 
  async onError(
    args: ProcessErrorArgs,
    handleError: (error: ServiceBusError) => Promise<void>
  ): Promise<void> {
    this.logger.error(
      `Error from source ${args.errorSource} occurred: `,
      args.error
    );
 
    Eif (isServiceBusError(args.error)) {
      switch (args.error.code) {
        case "MessagingEntityDisabled":
        case "MessagingEntityNotFound":
        case "UnauthorizedAccess":
          this.logger.error(
            `An unrecoverable error occurred. Stopping processing. ${args.error.code}`,
            args.error
          );
          await handleError(args.error);
          break;
        case "MessageLockLost":
          this.logger.error(`Message lock lost for message`, args.error);
          handleError(args.error);
          break;
        case "ServiceBusy":
          await delay(1000);
          handleError(args.error);
          break;
      }
    }
  }
}