- Flux<Policy> createTask() {
- return Flux.fromIterable(services.getAll()) //
- .filter(service -> service.isExpired()) //
- .doOnNext(service -> logger.info("Service is expired:" + service.getName()))
- .flatMap(service -> getAllPolicies(service)) //
- .flatMap(policy -> deletePolicy(policy));
+ private Flux<Policy> createTask() {
+ synchronized (services) {
+ return Flux.fromIterable(services.getAll()) //
+ .filter(Service::isExpired) //
+ .doOnNext(service -> logger.info("Service is expired: {}", service.getName())) //
+ .doOnNext(service -> services.remove(service.getName())) //
+ .flatMap(this::getAllPoliciesForService) //
+ .doOnNext(policies::remove) //
+ .flatMap(this::deletePolicyInRic);
+ }