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 | 2x 2x 138x 138x 138x 2x 84x 2x 12x 2x 11x 23x 2x 45x 45x 45x 84x 84x 9x 9x 9x 75x 75x 75x 63x 63x 12x 12x 7x 7x 5x 2x 2x 3x 3x 3x 5x 5x 5x | import { concatMap, delay, map, MonoTypeOperatorFunction, of, retry, Subject, take, tap } from 'rxjs';
import { IResilienceConfig, TTopicPropName } from '../../model/type/resilience.rx-operator.type';
export const getTopicSpecificBooleanOrDefault = (config: IResilienceConfig, propertyName: TTopicPropName): boolean => {
const topicConfig = config.topicToConfigDict[config.topic];
const isDefinedByTopicConfig = topicConfig && topicConfig[propertyName] !== undefined;
return isDefinedByTopicConfig ? (topicConfig[propertyName] as boolean) : config[propertyName];
};
const isRetryDisabled = (config: IResilienceConfig): boolean =>
getTopicSpecificBooleanOrDefault(config, 'disableRetry');
const waitForUserDecision = (config: IResilienceConfig): boolean =>
getTopicSpecificBooleanOrDefault(config, 'waitForUserDecision');
export const shouldLogResult = (config: IResilienceConfig): boolean =>
getTopicSpecificBooleanOrDefault(config, 'logResult');
export const shouldTrace = (config: IResilienceConfig): boolean => getTopicSpecificBooleanOrDefault(config, 'trace');
export const applyResilience = <T>(config: IResilienceConfig, uuid: string): MonoTypeOperatorFunction<T> => {
const userRetryOrCancelSubject = new Subject<boolean>();
let retryCount = 0;
return retry<T>({
delay: (error) =>
of(error).pipe(
concatMap((err) => {
if (isRetryDisabled(config) || !config.retryOnStatusCodeList.includes(err.status)) {
const topicSpecificMessageOrEmpty = config.topicToConfigDict[config.topic]?.failMessage || '';
config.onFail(config.topic, uuid, topicSpecificMessageOrEmpty);
throw err;
}
retryCount = retryCount + 1;
const retries = retryCount;
if (retryCount <= config.retryIntervalInMillisList.length) {
return of(err).pipe(
tap(() =>
config.onRequestRetry(
config.topic,
uuid,
retries,
config.retryIntervalInMillisList[retries - 1],
err.status,
),
),
delay(config.retryIntervalInMillisList[retryCount - 1]),
);
} else {
retryCount = 0;
if (waitForUserDecision(config)) {
config.onWaitingForUserDecision(
config.topic,
uuid,
retries,
err.status,
userRetryOrCancelSubject,
);
return userRetryOrCancelSubject.asObservable().pipe(
take(1),
map((shouldRetry) => {
if (shouldRetry) {
config.onRequestRetry(
config.topic,
uuid,
retries,
config.retryIntervalInMillisList[retries - 1],
err.status,
);
return of(err);
}
const topicSpecificMessageOrEmpty =
config.topicToConfigDict[config.topic]?.failMessage || '';
config.onFail(config.topic, uuid, topicSpecificMessageOrEmpty);
throw err;
}),
);
} else {
const topicSpecificMessageOrEmpty =
config.topicToConfigDict[config.topic]?.failMessage || '';
config.onFail(config.topic, uuid, topicSpecificMessageOrEmpty);
throw err;
}
}
}),
),
});
};
|