Больше года прошло с тех пор, как я познакомился с NgRx. На первый взгляд этот инструмент мне показался достаточно понятным. Но, чем больше я его использую, тем больше убеждаюсь, что это совсем не так. Тут я хотел бы отметить, что NgRx требует глубокого понимая принципов RxJS. Если в знании RxJS есть пробелы, придется не раз получить граблями по лбу.

Итак, я столкнулся со следующей проблемой – мой эффект срабатывал только один раз.

changeMode$ = createEffect(() => {
		return this.actions$.pipe(
			ofType(changeModeAction),
			map(action => action.mode),
			withLatestFrom(this.store.pipe(select(deviceStatesSelector))),
			map((responses: [string, DeviceState[]]) => {
				const devices = responses[1];
				const canChangeMode = devices?.length;
				if (!canChangeMode) {
					this.store.dispatch(presentToastAction({ message: 'No devices' }));
					throw new Error();
				}
				return responses;
			}),
			concatMap(([mode, devices]: [string, DeviceState[]]) => {
				return this.restService
					.changeMode(mode, devices)
					.pipe(map(() => changeModeSuccessAction()));
			}),
			catchError(() => of(changeModeFailureAction()))
		);
	});

Как вы видите из кода выше, в 11 строке я выбрасываю ошибку. Если девайсы отсутствуют, то режим менять нельзя. Я хочу прекратить работу эффекта на этом этапе и не отправлять запрос на сервер. Оператор catchError() перехватывает ошибку и возвращает действие, которое сигнализирует о неудачной попытке смены режима. После этого мой эффект перестает работать – я вызываю действие changeModeAction() второй, третий, четвертый раз, но больше ничего не происходит.

Я долго не мог понять, в чем дело. Перечитав документацию RxJS, через несколько часов мне удалось разобраться, как это работает. Оператор catchError() не просто перехватывает ошибку, он завершает исходный поток действий и мой эффект перестает работать.

Проблема была решена, когда я добавил в конце оператор repeat(). Он вернет Observable, который повторно подпишется на исходный поток, когда тот завершится.

changeMode$ = createEffect(() => {
		return this.actions$.pipe(
			ofType(changeModeAction),
			map(action => action.mode),
			withLatestFrom(this.store.pipe(select(deviceStatesSelector))),
			map((responses: [string, DeviceState[]]) => {
				const devices = responses[1];
				const canChangeMode = devices?.length;
				if (!canChangeMode) {
					this.store.dispatch(presentToastAction({ message: 'No devices' }));
					throw new Error();
				}
				return responses;
			}),
			concatMap(([mode, devices]: [string, DeviceState[]]) => {
				return this.restService
					.changeMode(mode, devices)
					.pipe(map(() => changeModeSuccessAction()));
			}),
			catchError(() => of(changeModeFailureAction())),
			repeat()
		);
	});

Есть еще один способ решить проблему, но он менее элегантный – можно завернуть все в concatMap(). Такие образом, мы переключимся с исходного потока действий на внутренний поток. Теперь ошибка отсутствия девайсов возникнет во внутреннем потоке и именно он будет завершен, а исходный поток продолжит свою работу.

changeMode$ = createEffect(() => {
		return this.actions$.pipe(
			ofType(changeModeAction),
			map(action => action.mode),
			withLatestFrom(this.store.pipe(select(deviceStatesSelector))),
			concatMap((responses: [string, DeviceState[]]) => {
				return of(responses).pipe(
					map((responses: [string, DeviceState[]]) => {
						const devices = responses[1];
						const canChangeMode = devices?.length;
						if (!canChangeMode) {
							this.store.dispatch(presentToastAction({ message: 'No devices' }));
							throw new Error();
						}
						return responses;
					}),
					concatMap(([mode, devices]: [string, DeviceState[]]) => {
						return this.restService
							.changeMode(mode, devices)
							.pipe(map(() => changeModeSuccessAction()));
					}),
					catchError(() => of(changeModeFailureAction()))
				);
			})
		);
	});

Надеюсь, статья была полезна для вас. Пишите комментарии – буду рад обратной связи. Успехов в реактивном программировании!