diff --git a/spring-cloud-app-broker-core/src/main/java/org/springframework/cloud/appbroker/service/WorkflowServiceInstanceBindingService.java b/spring-cloud-app-broker-core/src/main/java/org/springframework/cloud/appbroker/service/WorkflowServiceInstanceBindingService.java index 7fce998..353fffd 100644 --- a/spring-cloud-app-broker-core/src/main/java/org/springframework/cloud/appbroker/service/WorkflowServiceInstanceBindingService.java +++ b/spring-cloud-app-broker-core/src/main/java/org/springframework/cloud/appbroker/service/WorkflowServiceInstanceBindingService.java @@ -80,6 +80,7 @@ public class WorkflowServiceInstanceBindingService implements ServiceInstanceBin @Override public Mono createServiceInstanceBinding(CreateServiceInstanceBindingRequest request) { return invokeCreateResponseBuilders(request) + .publishOn(Schedulers.parallel()) .doOnNext(response -> create(request, response) .subscribe()); } @@ -132,7 +133,6 @@ public class WorkflowServiceInstanceBindingService implements ServiceInstanceBin CreateServiceInstanceBindingResponse response) { return stateRepository.saveState(request.getServiceInstanceId(), request.getBindingId(), OperationState.IN_PROGRESS, "create service instance binding started") - .publishOn(Schedulers.parallel()) .thenMany(invokeCreateWorkflows(request, response) .doOnRequest(l -> log.debug("Creating service instance binding")) .doOnComplete(() -> log.debug("Finished creating service instance binding")) @@ -168,6 +168,7 @@ public class WorkflowServiceInstanceBindingService implements ServiceInstanceBin @Override public Mono deleteServiceInstanceBinding(DeleteServiceInstanceBindingRequest request) { return invokeDeleteResponseBuilders(request) + .publishOn(Schedulers.parallel()) .doOnNext(response -> delete(request, response) .subscribe()); } @@ -188,7 +189,6 @@ public class WorkflowServiceInstanceBindingService implements ServiceInstanceBin DeleteServiceInstanceBindingResponse response) { return stateRepository.saveState(request.getServiceInstanceId(), request.getBindingId(), OperationState.IN_PROGRESS, "delete service instance binding started") - .publishOn(Schedulers.parallel()) .thenMany(invokeDeleteWorkflows(request, response) .doOnRequest(l -> log.debug("Deleting service instance binding")) .doOnComplete(() -> log.debug("Finished deleting service instance binding")) diff --git a/spring-cloud-app-broker-core/src/main/java/org/springframework/cloud/appbroker/service/WorkflowServiceInstanceService.java b/spring-cloud-app-broker-core/src/main/java/org/springframework/cloud/appbroker/service/WorkflowServiceInstanceService.java index d0d69f2..91a642d 100644 --- a/spring-cloud-app-broker-core/src/main/java/org/springframework/cloud/appbroker/service/WorkflowServiceInstanceService.java +++ b/spring-cloud-app-broker-core/src/main/java/org/springframework/cloud/appbroker/service/WorkflowServiceInstanceService.java @@ -77,6 +77,7 @@ public class WorkflowServiceInstanceService implements ServiceInstanceService { @Override public Mono createServiceInstance(CreateServiceInstanceRequest request) { return invokeCreateResponseBuilders(request) + .publishOn(Schedulers.parallel()) .doOnNext(response -> create(request, response) .subscribe()); } @@ -97,7 +98,6 @@ public class WorkflowServiceInstanceService implements ServiceInstanceService { return stateRepository.saveState(request.getServiceInstanceId(), OperationState.IN_PROGRESS, "create service instance started") - .publishOn(Schedulers.parallel()) .thenMany(invokeCreateWorkflows(request, response) .doOnRequest(l -> log.debug("Creating service instance")) .doOnComplete(() -> log.debug("Finished creating service instance")) @@ -121,6 +121,7 @@ public class WorkflowServiceInstanceService implements ServiceInstanceService { @Override public Mono deleteServiceInstance(DeleteServiceInstanceRequest request) { return invokeDeleteResponseBuilders(request) + .publishOn(Schedulers.parallel()) .doOnNext(response -> delete(request, response) .subscribe()); } @@ -140,7 +141,6 @@ public class WorkflowServiceInstanceService implements ServiceInstanceService { private Mono delete(DeleteServiceInstanceRequest request, DeleteServiceInstanceResponse response) { return stateRepository.saveState(request.getServiceInstanceId(), OperationState.IN_PROGRESS, "delete service instance started") - .publishOn(Schedulers.parallel()) .thenMany(invokeDeleteWorkflows(request, response) .doOnRequest(l -> log.debug("Deleting service instance")) .doOnComplete(() -> log.debug("Finished deleting service instance")) @@ -164,6 +164,7 @@ public class WorkflowServiceInstanceService implements ServiceInstanceService { @Override public Mono updateServiceInstance(UpdateServiceInstanceRequest request) { return invokeUpdateResponseBuilders(request) + .publishOn(Schedulers.parallel()) .doOnNext(response -> update(request, response) .subscribe()); } @@ -183,7 +184,6 @@ public class WorkflowServiceInstanceService implements ServiceInstanceService { private Mono update(UpdateServiceInstanceRequest request, UpdateServiceInstanceResponse response) { return stateRepository.saveState(request.getServiceInstanceId(), OperationState.IN_PROGRESS, "update service instance started") - .publishOn(Schedulers.parallel()) .thenMany(invokeUpdateWorkflows(request, response) .doOnRequest(l -> log.debug("Updating service instance")) .doOnComplete(() -> log.debug("Finished updating service instance")) diff --git a/spring-cloud-app-broker-core/src/test/java/org/springframework/cloud/appbroker/service/WorkflowServiceInstanceBindingServiceTest.java b/spring-cloud-app-broker-core/src/test/java/org/springframework/cloud/appbroker/service/WorkflowServiceInstanceBindingServiceTest.java index 72d8735..3fe020f 100644 --- a/spring-cloud-app-broker-core/src/test/java/org/springframework/cloud/appbroker/service/WorkflowServiceInstanceBindingServiceTest.java +++ b/spring-cloud-app-broker-core/src/test/java/org/springframework/cloud/appbroker/service/WorkflowServiceInstanceBindingServiceTest.java @@ -17,7 +17,6 @@ package org.springframework.cloud.appbroker.service; import java.sql.Timestamp; -import java.time.Duration; import java.time.Instant; import java.util.Arrays; @@ -55,6 +54,7 @@ import static org.mockito.ArgumentMatchers.anyString; import static org.mockito.ArgumentMatchers.eq; import static org.mockito.BDDMockito.given; import static org.mockito.Mockito.inOrder; +import static org.mockito.Mockito.verifyNoMoreInteractions; import static org.mockito.Mockito.when; @ExtendWith(MockitoExtension.class) @@ -160,7 +160,7 @@ class WorkflowServiceInstanceBindingServiceTest { } @Test - void createServiceInstanceAppBinding() throws InterruptedException { + void createServiceInstanceAppBinding() { when(stateRepository.saveState(anyString(), anyString(), any(OperationState.class), anyString())) .thenReturn(Mono.just(new ServiceInstanceState(OperationState.IN_PROGRESS, "create service instance binding started", new Timestamp(Instant.now().minusSeconds(60).toEpochMilli())))) @@ -197,21 +197,12 @@ class WorkflowServiceInstanceBindingServiceTest { given(createServiceInstanceAppBindingWorkflow2.accept(request)) .willReturn(Mono.just(true)); given(createServiceInstanceAppBindingWorkflow2.create(eq(request), eq(builtResponse))) - .willReturn(Mono.just(true) - .flatMap(value -> { - try { - Thread.sleep(500); - } catch (InterruptedException e) { - } - return higherOrderFlow.mono(); - })); + .willReturn(higherOrderFlow.mono()); given(createServiceInstanceAppBindingWorkflow2.buildResponse(eq(request), any(CreateServiceInstanceAppBindingResponseBuilder.class))) .willReturn(Mono.just(responseBuilder .credentials("foo", "bar") .operation("working2"))); - InOrder createOrder = inOrder(createServiceInstanceAppBindingWorkflow1, createServiceInstanceAppBindingWorkflow2); - StepVerifier.create(workflowServiceInstanceBindingService.createServiceInstanceBinding(request)) .assertNext(response -> { InOrder repoOrder = inOrder(stateRepository); @@ -221,10 +212,20 @@ class WorkflowServiceInstanceBindingServiceTest { .saveState(eq("foo-service"), eq("foo-binding"), eq(OperationState.SUCCEEDED), eq("create service instance binding completed")); repoOrder.verifyNoMoreInteractions(); + lowerOrderFlow.complete(); + lowerOrderFlow.assertWasNotRequested(); + + higherOrderFlow.complete(); + lowerOrderFlow.assertWasRequested(); + + InOrder createOrder = inOrder(createServiceInstanceAppBindingWorkflow1, createServiceInstanceAppBindingWorkflow2); createOrder.verify(createServiceInstanceAppBindingWorkflow2).buildResponse(eq(request), any(CreateServiceInstanceAppBindingResponseBuilder.class)); createOrder.verify(createServiceInstanceAppBindingWorkflow1).buildResponse(eq(request), any(CreateServiceInstanceAppBindingResponseBuilder.class)); + createOrder.verify(createServiceInstanceAppBindingWorkflow2).create(request, responseBuilder.build()); + createOrder.verify(createServiceInstanceAppBindingWorkflow1).create(request, responseBuilder.build()); + createOrder.verifyNoMoreInteractions(); CreateServiceInstanceAppBindingResponse r = (CreateServiceInstanceAppBindingResponse)response; @@ -233,20 +234,7 @@ class WorkflowServiceInstanceBindingServiceTest { assertThat(r.isAsync()).isTrue(); assertThat(r.getOperation()).isEqualTo("working2"); }) - .expectComplete() - .verifyThenAssertThat() - .tookLessThan(Duration.ofMillis(250)); - - lowerOrderFlow.complete(); - lowerOrderFlow.assertWasNotRequested(); - - higherOrderFlow.complete(); - Thread.sleep(600); - lowerOrderFlow.assertWasRequested(); - - createOrder.verify(createServiceInstanceAppBindingWorkflow2).create(request, responseBuilder.build()); - createOrder.verify(createServiceInstanceAppBindingWorkflow1).create(request, responseBuilder.build()); - createOrder.verifyNoMoreInteractions(); + .verifyComplete(); } @Test @@ -336,7 +324,7 @@ class WorkflowServiceInstanceBindingServiceTest { } @Test - void createServiceInstanceRouteBinding() throws InterruptedException { + void createServiceInstanceRouteBinding() { when(stateRepository.saveState(anyString(), anyString(), any(OperationState.class), anyString())) .thenReturn(Mono.just(new ServiceInstanceState(OperationState.IN_PROGRESS, "create service instance binding started", new Timestamp(Instant.now().minusSeconds(60).toEpochMilli())))) @@ -373,21 +361,12 @@ class WorkflowServiceInstanceBindingServiceTest { given(createServiceInstanceRouteBindingWorkflow2.accept(request)) .willReturn(Mono.just(true)); given(createServiceInstanceRouteBindingWorkflow2.create(eq(request), eq(builtResponse))) - .willReturn(Mono.just(true) - .flatMap(value -> { - try { - Thread.sleep(500); - } catch (InterruptedException e) { - } - return higherOrderFlow.mono(); - })); + .willReturn(higherOrderFlow.mono()); given(createServiceInstanceRouteBindingWorkflow2.buildResponse(eq(request), any(CreateServiceInstanceRouteBindingResponseBuilder.class))) .willReturn(Mono.just(responseBuilder .routeServiceUrl("foo-url") .operation("working2"))); - InOrder createOrder = inOrder(createServiceInstanceRouteBindingWorkflow1, createServiceInstanceRouteBindingWorkflow2); - StepVerifier.create(workflowServiceInstanceBindingService.createServiceInstanceBinding(request)) .assertNext(response -> { InOrder repoOrder = inOrder(stateRepository); @@ -396,11 +375,21 @@ class WorkflowServiceInstanceBindingServiceTest { repoOrder.verify(stateRepository) .saveState(eq("foo-service"), eq("foo-binding"), eq(OperationState.SUCCEEDED), eq("create service instance binding completed")); repoOrder.verifyNoMoreInteractions(); + + lowerOrderFlow.complete(); + lowerOrderFlow.assertWasNotRequested(); + higherOrderFlow.complete(); + lowerOrderFlow.assertWasRequested(); + + InOrder createOrder = inOrder(createServiceInstanceRouteBindingWorkflow1, createServiceInstanceRouteBindingWorkflow2); createOrder.verify(createServiceInstanceRouteBindingWorkflow2).buildResponse(eq(request), any(CreateServiceInstanceRouteBindingResponseBuilder.class)); createOrder.verify(createServiceInstanceRouteBindingWorkflow1).buildResponse(eq(request), any(CreateServiceInstanceRouteBindingResponseBuilder.class)); + createOrder.verify(createServiceInstanceRouteBindingWorkflow2).create(request, responseBuilder.build()); + createOrder.verify(createServiceInstanceRouteBindingWorkflow1).create(request, responseBuilder.build()); + createOrder.verifyNoMoreInteractions(); assertThat(response).isNotNull(); assertThat(response).isInstanceOf(CreateServiceInstanceRouteBindingResponse.class); @@ -408,24 +397,11 @@ class WorkflowServiceInstanceBindingServiceTest { assertThat(response.isAsync()).isTrue(); assertThat(response.getOperation()).isEqualTo("working2"); }) - .expectComplete() - .verifyThenAssertThat() - .tookLessThan(Duration.ofMillis(250)); - - lowerOrderFlow.complete(); - lowerOrderFlow.assertWasNotRequested(); - - higherOrderFlow.complete(); - Thread.sleep(600); - lowerOrderFlow.assertWasRequested(); - - createOrder.verify(createServiceInstanceRouteBindingWorkflow2).create(request, responseBuilder.build()); - createOrder.verify(createServiceInstanceRouteBindingWorkflow1).create(request, responseBuilder.build()); - createOrder.verifyNoMoreInteractions(); + .verifyComplete(); } @Test - void createServiceInstanceAppBindingWithAsyncError() throws InterruptedException { + void createServiceInstanceAppBindingWithAsyncError() { when(stateRepository.saveState(anyString(), anyString(), any(OperationState.class), anyString())) .thenReturn(Mono.just(new ServiceInstanceState(OperationState.IN_PROGRESS, "create service instance binding started", new Timestamp(Instant.now().minusSeconds(60).toEpochMilli())))) @@ -456,36 +432,32 @@ class WorkflowServiceInstanceBindingServiceTest { given(createServiceInstanceAppBindingWorkflow2.buildResponse(eq(request), any(CreateServiceInstanceAppBindingResponseBuilder.class))) .willReturn(Mono.just(responseBuilder)); - InOrder repoOrder = inOrder(stateRepository); - InOrder createOrder = inOrder(createServiceInstanceAppBindingWorkflow1, createServiceInstanceAppBindingWorkflow2); - StepVerifier.create(workflowServiceInstanceBindingService.createServiceInstanceBinding(request)) .assertNext(response -> { + InOrder repoOrder = inOrder(stateRepository); repoOrder.verify(stateRepository) .saveState(eq("foo-service"), eq("foo-binding"), eq(OperationState.IN_PROGRESS), eq("create service instance binding started")); + repoOrder.verify(stateRepository) + .saveState(eq("foo-service"), eq("foo-binding"), eq(OperationState.FAILED), eq("create foo error")); + repoOrder.verifyNoMoreInteractions(); + InOrder createOrder = inOrder(createServiceInstanceAppBindingWorkflow1, createServiceInstanceAppBindingWorkflow2); createOrder.verify(createServiceInstanceAppBindingWorkflow2).buildResponse(eq(request), any(CreateServiceInstanceAppBindingResponseBuilder.class)); createOrder.verify(createServiceInstanceAppBindingWorkflow1).buildResponse(eq(request), any(CreateServiceInstanceAppBindingResponseBuilder.class)); + createOrder.verify(createServiceInstanceAppBindingWorkflow2).create(request, responseBuilder.build()); + createOrder.verify(createServiceInstanceAppBindingWorkflow1).create(request, responseBuilder.build()); + createOrder.verifyNoMoreInteractions(); assertThat(response).isNotNull(); assertThat(response).isInstanceOf(CreateServiceInstanceAppBindingResponse.class); }) .verifyComplete(); - - Thread.sleep(50); - - repoOrder.verify(stateRepository).saveState(eq("foo-service"), eq("foo-binding"), eq(OperationState.FAILED), eq("create foo error")); - repoOrder.verifyNoMoreInteractions(); - - createOrder.verify(createServiceInstanceAppBindingWorkflow2).create(request, responseBuilder.build()); - createOrder.verify(createServiceInstanceAppBindingWorkflow1).create(request, responseBuilder.build()); - createOrder.verifyNoMoreInteractions(); } @Test - void createServiceInstanceRouteBindingWithAsyncError() throws InterruptedException { + void createServiceInstanceRouteBindingWithAsyncError() { when(stateRepository.saveState(anyString(), anyString(), any(OperationState.class), anyString())) .thenReturn(Mono.just(new ServiceInstanceState(OperationState.IN_PROGRESS, "create service instance binding started", new Timestamp(Instant.now().minusSeconds(60).toEpochMilli())))) @@ -516,32 +488,28 @@ class WorkflowServiceInstanceBindingServiceTest { given(createServiceInstanceRouteBindingWorkflow2.buildResponse(eq(request), any(CreateServiceInstanceRouteBindingResponseBuilder.class))) .willReturn(Mono.just(responseBuilder)); - InOrder repoOrder = inOrder(stateRepository); - InOrder createOrder = inOrder(createServiceInstanceRouteBindingWorkflow1, createServiceInstanceRouteBindingWorkflow2); - StepVerifier.create(workflowServiceInstanceBindingService.createServiceInstanceBinding(request)) .assertNext(response -> { + InOrder repoOrder = inOrder(stateRepository); repoOrder.verify(stateRepository) .saveState(eq("foo-service"), eq("foo-binding"), eq(OperationState.IN_PROGRESS), eq("create service instance binding started")); + repoOrder.verify(stateRepository) + .saveState(eq("foo-service"), eq("foo-binding"), eq(OperationState.FAILED), eq("create foo error")); + repoOrder.verifyNoMoreInteractions(); + InOrder createOrder = inOrder(createServiceInstanceRouteBindingWorkflow1, createServiceInstanceRouteBindingWorkflow2); createOrder.verify(createServiceInstanceRouteBindingWorkflow2).buildResponse(eq(request), any(CreateServiceInstanceRouteBindingResponseBuilder.class)); createOrder.verify(createServiceInstanceRouteBindingWorkflow1).buildResponse(eq(request), any(CreateServiceInstanceRouteBindingResponseBuilder.class)); + createOrder.verify(createServiceInstanceRouteBindingWorkflow2).create(request, responseBuilder.build()); + createOrder.verify(createServiceInstanceRouteBindingWorkflow1).create(request, responseBuilder.build()); + createOrder.verifyNoMoreInteractions(); assertThat(response).isNotNull(); assertThat(response).isInstanceOf(CreateServiceInstanceRouteBindingResponse.class); }) .verifyComplete(); - - Thread.sleep(50); - - repoOrder.verify(stateRepository).saveState(eq("foo-service"), eq("foo-binding"), eq(OperationState.FAILED), eq("create foo error")); - repoOrder.verifyNoMoreInteractions(); - - createOrder.verify(createServiceInstanceRouteBindingWorkflow2).create(request, responseBuilder.build()); - createOrder.verify(createServiceInstanceRouteBindingWorkflow1).create(request, responseBuilder.build()); - createOrder.verifyNoMoreInteractions(); } @Test @@ -601,7 +569,7 @@ class WorkflowServiceInstanceBindingServiceTest { } @Test - void createServiceInstanceAppBindingWithNoAcceptsDoesNothing() throws InterruptedException { + void createServiceInstanceAppBindingWithNoAcceptsDoesNothing() { when(stateRepository.saveState(anyString(), anyString(), any(OperationState.class), anyString())) .thenReturn(Mono.just(new ServiceInstanceState(OperationState.IN_PROGRESS, "create service instance started", new Timestamp(Instant.now().minusSeconds(60).toEpochMilli())))) @@ -622,8 +590,6 @@ class WorkflowServiceInstanceBindingServiceTest { given(createServiceInstanceAppBindingWorkflow2.accept(request)) .willReturn(Mono.just(false)); - InOrder createOrder = inOrder(createServiceInstanceAppBindingWorkflow1, createServiceInstanceAppBindingWorkflow2); - StepVerifier.create(workflowServiceInstanceBindingService.createServiceInstanceBinding(request)) .assertNext(response -> { InOrder repoOrder = inOrder(stateRepository); @@ -633,23 +599,16 @@ class WorkflowServiceInstanceBindingServiceTest { .saveState(eq("foo-service"), eq("foo-binding"), eq(OperationState.SUCCEEDED), eq("create service instance binding completed")); repoOrder.verifyNoMoreInteractions(); - createOrder.verify(createServiceInstanceAppBindingWorkflow2).accept(eq(request)); - createOrder.verify(createServiceInstanceAppBindingWorkflow1).accept(eq(request)); + verifyNoMoreInteractions(createServiceInstanceAppBindingWorkflow1, createServiceInstanceAppBindingWorkflow2); assertThat(response).isNotNull(); assertThat(response).isInstanceOf(CreateServiceInstanceAppBindingResponse.class); }) .verifyComplete(); - - Thread.sleep(50); - - createOrder.verify(createServiceInstanceAppBindingWorkflow2).accept(eq(request)); - createOrder.verify(createServiceInstanceAppBindingWorkflow1).accept(eq(request)); - createOrder.verifyNoMoreInteractions(); } @Test - void createServiceInstanceRouteBindingWithNoAcceptsDoesNothing() throws InterruptedException { + void createServiceInstanceRouteBindingWithNoAcceptsDoesNothing() { when(stateRepository.saveState(anyString(), anyString(), any(OperationState.class), anyString())) .thenReturn(Mono.just(new ServiceInstanceState(OperationState.IN_PROGRESS, "create service instance started", new Timestamp(Instant.now().minusSeconds(60).toEpochMilli())))) @@ -670,8 +629,6 @@ class WorkflowServiceInstanceBindingServiceTest { given(createServiceInstanceRouteBindingWorkflow2.accept(request)) .willReturn(Mono.just(false)); - InOrder createOrder = inOrder(createServiceInstanceRouteBindingWorkflow1, createServiceInstanceRouteBindingWorkflow2); - StepVerifier.create(workflowServiceInstanceBindingService.createServiceInstanceBinding(request)) .assertNext(response -> { InOrder repoOrder = inOrder(stateRepository); @@ -681,19 +638,12 @@ class WorkflowServiceInstanceBindingServiceTest { .saveState(eq("foo-service"), eq("foo-binding"), eq(OperationState.SUCCEEDED), eq("create service instance binding completed")); repoOrder.verifyNoMoreInteractions(); - createOrder.verify(createServiceInstanceRouteBindingWorkflow2).accept(eq(request)); - createOrder.verify(createServiceInstanceRouteBindingWorkflow1).accept(eq(request)); + verifyNoMoreInteractions(createServiceInstanceRouteBindingWorkflow1, createServiceInstanceRouteBindingWorkflow2); assertThat(response).isNotNull(); assertThat(response).isInstanceOf(CreateServiceInstanceRouteBindingResponse.class); }) .verifyComplete(); - - Thread.sleep(50); - - createOrder.verify(createServiceInstanceRouteBindingWorkflow2).accept(eq(request)); - createOrder.verify(createServiceInstanceRouteBindingWorkflow1).accept(eq(request)); - createOrder.verifyNoMoreInteractions(); } @Test @@ -728,7 +678,7 @@ class WorkflowServiceInstanceBindingServiceTest { } @Test - void deleteServiceInstanceBinding() throws InterruptedException { + void deleteServiceInstanceBinding() { when(stateRepository.saveState(anyString(), anyString(), any(OperationState.class), anyString())) .thenReturn(Mono.just(new ServiceInstanceState(OperationState.IN_PROGRESS, "delete service instance binding started", new Timestamp(Instant.now().minusSeconds(60).toEpochMilli())))) @@ -752,7 +702,7 @@ class WorkflowServiceInstanceBindingServiceTest { given(deleteServiceInstanceBindingWorkflow1.accept(request)) .willReturn(Mono.just(true)); given(deleteServiceInstanceBindingWorkflow1.delete(eq(request), eq(builtResponse))) - .willReturn(lowerOrderFlow.mono().then()); + .willReturn(lowerOrderFlow.mono()); given(deleteServiceInstanceBindingWorkflow1.buildResponse(eq(request), any(DeleteServiceInstanceBindingResponseBuilder.class))) .willReturn(Mono.just(responseBuilder .async(true) @@ -761,20 +711,11 @@ class WorkflowServiceInstanceBindingServiceTest { given(deleteServiceInstanceBindingWorkflow2.accept(request)) .willReturn(Mono.just(true)); given(deleteServiceInstanceBindingWorkflow2.delete(eq(request), eq(builtResponse))) - .willReturn(Mono.just(true) - .flatMap(value -> { - try { - Thread.sleep(500); - } catch (InterruptedException e) { - } - return higherOrderFlow.mono(); - })); + .willReturn(higherOrderFlow.mono()); given(deleteServiceInstanceBindingWorkflow2.buildResponse(eq(request), any(DeleteServiceInstanceBindingResponseBuilder.class))) .willReturn(Mono.just(responseBuilder .operation("working2"))); - InOrder deleteOrder = inOrder(deleteServiceInstanceBindingWorkflow1, deleteServiceInstanceBindingWorkflow2); - StepVerifier.create(workflowServiceInstanceBindingService.deleteServiceInstanceBinding(request)) .assertNext(response -> { InOrder repoOrder = inOrder(stateRepository); @@ -784,33 +725,30 @@ class WorkflowServiceInstanceBindingServiceTest { .saveState(eq("foo-service"), eq("foo-binding"), eq(OperationState.SUCCEEDED), eq("delete service instance binding completed")); repoOrder.verifyNoMoreInteractions(); + lowerOrderFlow.complete(); + lowerOrderFlow.assertWasNotRequested(); + + higherOrderFlow.complete(); + lowerOrderFlow.assertWasRequested(); + + InOrder deleteOrder = inOrder(deleteServiceInstanceBindingWorkflow1, deleteServiceInstanceBindingWorkflow2); deleteOrder.verify(deleteServiceInstanceBindingWorkflow2).buildResponse(eq(request), any(DeleteServiceInstanceBindingResponseBuilder.class)); deleteOrder.verify(deleteServiceInstanceBindingWorkflow1).buildResponse(eq(request), any(DeleteServiceInstanceBindingResponseBuilder.class)); + deleteOrder.verify(deleteServiceInstanceBindingWorkflow2).delete(request, responseBuilder.build()); + deleteOrder.verify(deleteServiceInstanceBindingWorkflow1).delete(request, responseBuilder.build()); + deleteOrder.verifyNoMoreInteractions(); assertThat(response).isNotNull(); assertThat(response.isAsync()).isTrue(); assertThat(response.getOperation()).isEqualTo("working2"); }) - .expectComplete() - .verifyThenAssertThat() - .tookLessThan(Duration.ofMillis(250)); - - lowerOrderFlow.complete(); - lowerOrderFlow.assertWasNotRequested(); - - higherOrderFlow.complete(); - Thread.sleep(600); - lowerOrderFlow.assertWasRequested(); - - deleteOrder.verify(deleteServiceInstanceBindingWorkflow2).delete(request, responseBuilder.build()); - deleteOrder.verify(deleteServiceInstanceBindingWorkflow1).delete(request, responseBuilder.build()); - deleteOrder.verifyNoMoreInteractions(); + .verifyComplete(); } @Test - void deleteServiceInstanceBindingWithAsyncError() throws InterruptedException { + void deleteServiceInstanceBindingWithAsyncError() { when(stateRepository.saveState(anyString(), anyString(), any(OperationState.class), anyString())) .thenReturn(Mono.just(new ServiceInstanceState(OperationState.IN_PROGRESS, "delete service instance started", new Timestamp(Instant.now().minusSeconds(60).toEpochMilli())))) @@ -838,31 +776,27 @@ class WorkflowServiceInstanceBindingServiceTest { given(deleteServiceInstanceBindingWorkflow2.buildResponse(eq(request), any(DeleteServiceInstanceBindingResponseBuilder.class))) .willReturn(Mono.just(responseBuilder)); - InOrder repoOrder = inOrder(stateRepository); - InOrder deleteOrder = inOrder(deleteServiceInstanceBindingWorkflow1, deleteServiceInstanceBindingWorkflow2); - StepVerifier.create(workflowServiceInstanceBindingService.deleteServiceInstanceBinding(request)) .assertNext(response -> { + InOrder repoOrder = inOrder(stateRepository); repoOrder.verify(stateRepository) .saveState(eq("foo-service"), eq("foo-binding"), eq(OperationState.IN_PROGRESS), eq("delete service instance binding started")); + repoOrder.verify(stateRepository) + .saveState(eq("foo-service"), eq("foo-binding"), eq(OperationState.FAILED), eq("delete foo binding error")); + repoOrder.verifyNoMoreInteractions(); + InOrder deleteOrder = inOrder(deleteServiceInstanceBindingWorkflow1, deleteServiceInstanceBindingWorkflow2); deleteOrder.verify(deleteServiceInstanceBindingWorkflow2).buildResponse(eq(request), any(DeleteServiceInstanceBindingResponseBuilder.class)); deleteOrder.verify(deleteServiceInstanceBindingWorkflow1).buildResponse(eq(request), any(DeleteServiceInstanceBindingResponseBuilder.class)); + deleteOrder.verify(deleteServiceInstanceBindingWorkflow2).delete(request, responseBuilder.build()); + deleteOrder.verify(deleteServiceInstanceBindingWorkflow1).delete(request, responseBuilder.build()); + deleteOrder.verifyNoMoreInteractions(); assertThat(response).isNotNull(); }) .verifyComplete(); - - Thread.sleep(50); - - repoOrder.verify(stateRepository).saveState(eq("foo-service"), eq("foo-binding"), eq(OperationState.FAILED), eq("delete foo binding error")); - repoOrder.verifyNoMoreInteractions(); - - deleteOrder.verify(deleteServiceInstanceBindingWorkflow2).delete(request, responseBuilder.build()); - deleteOrder.verify(deleteServiceInstanceBindingWorkflow1).delete(request, responseBuilder.build()); - deleteOrder.verifyNoMoreInteractions(); } @Test @@ -892,7 +826,7 @@ class WorkflowServiceInstanceBindingServiceTest { } @Test - void deleteServiceInstanceBindingWithNoAcceptsDoesNothing() throws InterruptedException { + void deleteServiceInstanceBindingWithNoAcceptsDoesNothing() { when(stateRepository.saveState(anyString(), anyString(), any(OperationState.class), anyString())) .thenReturn(Mono.just(new ServiceInstanceState(OperationState.IN_PROGRESS, "delete service instance binding started", new Timestamp(Instant.now().minusSeconds(60).toEpochMilli())))) @@ -910,8 +844,6 @@ class WorkflowServiceInstanceBindingServiceTest { given(deleteServiceInstanceBindingWorkflow2.accept(request)) .willReturn(Mono.just(false)); - InOrder deleteOrder = inOrder(deleteServiceInstanceBindingWorkflow1, deleteServiceInstanceBindingWorkflow2); - StepVerifier.create(workflowServiceInstanceBindingService.deleteServiceInstanceBinding(request)) .assertNext(response -> { InOrder repoOrder = inOrder(stateRepository); @@ -921,18 +853,11 @@ class WorkflowServiceInstanceBindingServiceTest { .saveState(eq("foo-service"), eq("foo-binding"), eq(OperationState.SUCCEEDED), eq("delete service instance binding completed")); repoOrder.verifyNoMoreInteractions(); - deleteOrder.verify(deleteServiceInstanceBindingWorkflow2).accept(eq(request)); - deleteOrder.verify(deleteServiceInstanceBindingWorkflow1).accept(eq(request)); + verifyNoMoreInteractions(deleteServiceInstanceBindingWorkflow1, deleteServiceInstanceBindingWorkflow2); assertThat(response).isNotNull(); }) .verifyComplete(); - - Thread.sleep(50); - - deleteOrder.verify(deleteServiceInstanceBindingWorkflow2).accept(eq(request)); - deleteOrder.verify(deleteServiceInstanceBindingWorkflow1).accept(eq(request)); - deleteOrder.verifyNoMoreInteractions(); } @Order(Ordered.HIGHEST_PRECEDENCE) diff --git a/spring-cloud-app-broker-core/src/test/java/org/springframework/cloud/appbroker/service/WorkflowServiceInstanceServiceTest.java b/spring-cloud-app-broker-core/src/test/java/org/springframework/cloud/appbroker/service/WorkflowServiceInstanceServiceTest.java index 92c33cb..6a881ec 100644 --- a/spring-cloud-app-broker-core/src/test/java/org/springframework/cloud/appbroker/service/WorkflowServiceInstanceServiceTest.java +++ b/spring-cloud-app-broker-core/src/test/java/org/springframework/cloud/appbroker/service/WorkflowServiceInstanceServiceTest.java @@ -17,7 +17,6 @@ package org.springframework.cloud.appbroker.service; import java.sql.Timestamp; -import java.time.Duration; import java.time.Instant; import java.util.Arrays; @@ -53,6 +52,7 @@ import static org.mockito.ArgumentMatchers.anyString; import static org.mockito.ArgumentMatchers.eq; import static org.mockito.BDDMockito.given; import static org.mockito.Mockito.inOrder; +import static org.mockito.Mockito.verifyNoMoreInteractions; import static org.mockito.Mockito.when; @ExtendWith(MockitoExtension.class) @@ -90,7 +90,7 @@ class WorkflowServiceInstanceServiceTest { } @Test - void createServiceInstance() throws InterruptedException { + void createServiceInstance() { when(serviceInstanceStateRepository.saveState(anyString(), any(OperationState.class), anyString())) .thenReturn(Mono.just(new ServiceInstanceState(OperationState.IN_PROGRESS, "create service instance started", new Timestamp(Instant.now().minusSeconds(60).toEpochMilli())))) @@ -123,21 +123,12 @@ class WorkflowServiceInstanceServiceTest { given(createServiceInstanceWorkflow2.accept(request)) .willReturn(Mono.just(true)); given(createServiceInstanceWorkflow2.create(eq(request), eq(builtResponse))) - .willReturn(Mono.just(true) - .flatMap(value -> { - try { - Thread.sleep(500); - } catch (InterruptedException e) { - } - return higherOrderFlow.mono(); - })); + .willReturn(higherOrderFlow.mono()); given(createServiceInstanceWorkflow2.buildResponse(eq(request), any(CreateServiceInstanceResponseBuilder.class))) .willReturn(Mono.just(responseBuilder .dashboardUrl("https://dashboard.example.com") .operation("working2"))); - InOrder createOrder = inOrder(createServiceInstanceWorkflow1, createServiceInstanceWorkflow2); - StepVerifier.create(workflowServiceInstanceService.createServiceInstance(request)) .assertNext(response -> { InOrder repoOrder = inOrder(serviceInstanceStateRepository); @@ -147,34 +138,31 @@ class WorkflowServiceInstanceServiceTest { .saveState(eq("foo"), eq(OperationState.SUCCEEDED), eq("create service instance completed")); repoOrder.verifyNoMoreInteractions(); + lowerOrderFlow.complete(); + lowerOrderFlow.assertWasNotRequested(); + + higherOrderFlow.complete(); + lowerOrderFlow.assertWasRequested(); + + InOrder createOrder = inOrder(createServiceInstanceWorkflow1, createServiceInstanceWorkflow2); createOrder.verify(createServiceInstanceWorkflow2).buildResponse(eq(request), any(CreateServiceInstanceResponseBuilder.class)); createOrder.verify(createServiceInstanceWorkflow1).buildResponse(eq(request), any(CreateServiceInstanceResponseBuilder.class)); + createOrder.verify(createServiceInstanceWorkflow2).create(request, responseBuilder.build()); + createOrder.verify(createServiceInstanceWorkflow1).create(request, responseBuilder.build()); + createOrder.verifyNoMoreInteractions(); assertThat(response).isNotNull(); assertThat(response.isAsync()).isTrue(); assertThat(response.getDashboardUrl()).isEqualTo("https://dashboard.example.com"); assertThat(response.getOperation()).isEqualTo("working2"); }) - .expectComplete() - .verifyThenAssertThat() - .tookLessThan(Duration.ofMillis(250)); - - lowerOrderFlow.complete(); - lowerOrderFlow.assertWasNotRequested(); - - higherOrderFlow.complete(); - Thread.sleep(600); - lowerOrderFlow.assertWasRequested(); - - createOrder.verify(createServiceInstanceWorkflow2).create(request, responseBuilder.build()); - createOrder.verify(createServiceInstanceWorkflow1).create(request, responseBuilder.build()); - createOrder.verifyNoMoreInteractions(); + .verifyComplete(); } @Test - void createServiceInstanceWithAsyncError() throws InterruptedException { + void createServiceInstanceWithAsyncError() { when(serviceInstanceStateRepository.saveState(anyString(), any(OperationState.class), anyString())) .thenReturn(Mono.just(new ServiceInstanceState(OperationState.IN_PROGRESS, "create service instance started", new Timestamp(Instant.now().minusSeconds(60).toEpochMilli())))) @@ -201,31 +189,27 @@ class WorkflowServiceInstanceServiceTest { given(createServiceInstanceWorkflow2.buildResponse(eq(request), any(CreateServiceInstanceResponseBuilder.class))) .willReturn(Mono.just(responseBuilder)); - InOrder repoOrder = inOrder(serviceInstanceStateRepository); - InOrder createOrder = inOrder(createServiceInstanceWorkflow1, createServiceInstanceWorkflow2); - StepVerifier.create(workflowServiceInstanceService.createServiceInstance(request)) .assertNext(response -> { + InOrder repoOrder = inOrder(serviceInstanceStateRepository); repoOrder.verify(serviceInstanceStateRepository) .saveState(eq("foo"), eq(OperationState.IN_PROGRESS), eq("create service instance started")); + repoOrder.verify(serviceInstanceStateRepository) + .saveState(eq("foo"), eq(OperationState.FAILED), eq("create foo error")); + repoOrder.verifyNoMoreInteractions(); + InOrder createOrder = inOrder(createServiceInstanceWorkflow1, createServiceInstanceWorkflow2); createOrder.verify(createServiceInstanceWorkflow2).buildResponse(eq(request), any(CreateServiceInstanceResponseBuilder.class)); createOrder.verify(createServiceInstanceWorkflow1).buildResponse(eq(request), any(CreateServiceInstanceResponseBuilder.class)); + createOrder.verify(createServiceInstanceWorkflow2).create(request, responseBuilder.build()); + createOrder.verify(createServiceInstanceWorkflow1).create(request, responseBuilder.build()); + createOrder.verifyNoMoreInteractions(); assertThat(response).isNotNull(); }) .verifyComplete(); - - Thread.sleep(50); - - repoOrder.verify(serviceInstanceStateRepository).saveState(eq("foo"), eq(OperationState.FAILED), eq("create foo error")); - repoOrder.verifyNoMoreInteractions(); - - createOrder.verify(createServiceInstanceWorkflow2).create(request, responseBuilder.build()); - createOrder.verify(createServiceInstanceWorkflow1).create(request, responseBuilder.build()); - createOrder.verifyNoMoreInteractions(); } @Test @@ -254,7 +238,7 @@ class WorkflowServiceInstanceServiceTest { } @Test - void createServiceInstanceWithNoAcceptsDoesNothing() throws InterruptedException { + void createServiceInstanceWithNoAcceptsDoesNothing() { when(serviceInstanceStateRepository.saveState(anyString(), any(OperationState.class), anyString())) .thenReturn(Mono.just(new ServiceInstanceState(OperationState.IN_PROGRESS, "create service instance started", new Timestamp(Instant.now().minusSeconds(60).toEpochMilli())))) @@ -271,8 +255,6 @@ class WorkflowServiceInstanceServiceTest { given(createServiceInstanceWorkflow2.accept(request)) .willReturn(Mono.just(false)); - InOrder createOrder = inOrder(createServiceInstanceWorkflow1, createServiceInstanceWorkflow2); - StepVerifier.create(workflowServiceInstanceService.createServiceInstance(request)) .assertNext(response -> { InOrder repoOrder = inOrder(serviceInstanceStateRepository); @@ -282,22 +264,15 @@ class WorkflowServiceInstanceServiceTest { .saveState(eq("foo"), eq(OperationState.SUCCEEDED), eq("create service instance completed")); repoOrder.verifyNoMoreInteractions(); - createOrder.verify(createServiceInstanceWorkflow2).accept(eq(request)); - createOrder.verify(createServiceInstanceWorkflow1).accept(eq(request)); + verifyNoMoreInteractions(createServiceInstanceWorkflow1, createServiceInstanceWorkflow2); assertThat(response).isNotNull(); }) .verifyComplete(); - - Thread.sleep(50); - - createOrder.verify(createServiceInstanceWorkflow2).accept(eq(request)); - createOrder.verify(createServiceInstanceWorkflow1).accept(eq(request)); - createOrder.verifyNoMoreInteractions(); } @Test - void deleteServiceInstance() throws InterruptedException { + void deleteServiceInstance() { when(serviceInstanceStateRepository.saveState(anyString(), any(OperationState.class), anyString())) .thenReturn(Mono.just(new ServiceInstanceState(OperationState.IN_PROGRESS, "delete service instance started", new Timestamp(Instant.now().minusSeconds(60).toEpochMilli())))) @@ -332,16 +307,7 @@ class WorkflowServiceInstanceServiceTest { .willReturn(Mono.just(responseBuilder .operation("working2"))); given(deleteServiceInstanceWorkflow2.delete(eq(request), eq(builtResponse))) - .willReturn(Mono.just(true) - .flatMap(value -> { - try { - Thread.sleep(500); - } catch (InterruptedException e) { - } - return higherOrderFlow.mono(); - })); - - InOrder deleteOrder = inOrder(deleteServiceInstanceWorkflow1, deleteServiceInstanceWorkflow2); + .willReturn(higherOrderFlow.mono()); StepVerifier.create(workflowServiceInstanceService.deleteServiceInstance(request)) .assertNext(response -> { @@ -352,33 +318,30 @@ class WorkflowServiceInstanceServiceTest { .saveState(eq("foo"), eq(OperationState.SUCCEEDED), eq("delete service instance completed")); repoOrder.verifyNoMoreInteractions(); + lowerOrderFlow.complete(); + lowerOrderFlow.assertWasNotRequested(); + + higherOrderFlow.complete(); + lowerOrderFlow.assertWasRequested(); + + InOrder deleteOrder = inOrder(deleteServiceInstanceWorkflow1, deleteServiceInstanceWorkflow2); deleteOrder.verify(deleteServiceInstanceWorkflow2).buildResponse(eq(request), any(DeleteServiceInstanceResponseBuilder.class)); deleteOrder.verify(deleteServiceInstanceWorkflow1).buildResponse(eq(request), any(DeleteServiceInstanceResponseBuilder.class)); + deleteOrder.verify(deleteServiceInstanceWorkflow2).delete(request, responseBuilder.build()); + deleteOrder.verify(deleteServiceInstanceWorkflow1).delete(request, responseBuilder.build()); + deleteOrder.verifyNoMoreInteractions(); assertThat(response).isNotNull(); assertThat(response.isAsync()).isTrue(); assertThat(response.getOperation()).isEqualTo("working2"); }) - .expectComplete() - .verifyThenAssertThat() - .tookLessThan(Duration.ofMillis(250)); - - lowerOrderFlow.complete(); - lowerOrderFlow.assertWasNotRequested(); - - higherOrderFlow.complete(); - Thread.sleep(600); - lowerOrderFlow.assertWasRequested(); - - deleteOrder.verify(deleteServiceInstanceWorkflow2).delete(request, responseBuilder.build()); - deleteOrder.verify(deleteServiceInstanceWorkflow1).delete(request, responseBuilder.build()); - deleteOrder.verifyNoMoreInteractions(); + .verifyComplete(); } @Test - void deleteServiceInstanceWithAsyncError() throws InterruptedException { + void deleteServiceInstanceWithAsyncError() { when(serviceInstanceStateRepository.saveState(anyString(), any(OperationState.class), anyString())) .thenReturn(Mono.just(new ServiceInstanceState(OperationState.IN_PROGRESS, "delete service instance started", new Timestamp(Instant.now().minusSeconds(60).toEpochMilli())))) @@ -405,31 +368,27 @@ class WorkflowServiceInstanceServiceTest { given(deleteServiceInstanceWorkflow2.buildResponse(eq(request), any(DeleteServiceInstanceResponseBuilder.class))) .willReturn(Mono.just(responseBuilder)); - InOrder repoOrder = inOrder(serviceInstanceStateRepository); - InOrder deleteOrder = inOrder(deleteServiceInstanceWorkflow1, deleteServiceInstanceWorkflow2); - StepVerifier.create(workflowServiceInstanceService.deleteServiceInstance(request)) .assertNext(response -> { + InOrder repoOrder = inOrder(serviceInstanceStateRepository); repoOrder.verify(serviceInstanceStateRepository) .saveState(eq("foo"), eq(OperationState.IN_PROGRESS), eq("delete service instance started")); + repoOrder.verify(serviceInstanceStateRepository) + .saveState(eq("foo"), eq(OperationState.FAILED), eq("delete foo error")); + repoOrder.verifyNoMoreInteractions(); + InOrder deleteOrder = inOrder(deleteServiceInstanceWorkflow1, deleteServiceInstanceWorkflow2); deleteOrder.verify(deleteServiceInstanceWorkflow2).buildResponse(eq(request), any(DeleteServiceInstanceResponseBuilder.class)); deleteOrder.verify(deleteServiceInstanceWorkflow1).buildResponse(eq(request), any(DeleteServiceInstanceResponseBuilder.class)); + deleteOrder.verify(deleteServiceInstanceWorkflow2).delete(request, responseBuilder.build()); + deleteOrder.verify(deleteServiceInstanceWorkflow1).delete(request, responseBuilder.build()); + deleteOrder.verifyNoMoreInteractions(); assertThat(response).isNotNull(); }) .verifyComplete(); - - Thread.sleep(50); - - repoOrder.verify(serviceInstanceStateRepository).saveState(eq("foo"), eq(OperationState.FAILED), eq("delete foo error")); - repoOrder.verifyNoMoreInteractions(); - - deleteOrder.verify(deleteServiceInstanceWorkflow2).delete(request, responseBuilder.build()); - deleteOrder.verify(deleteServiceInstanceWorkflow1).delete(request, responseBuilder.build()); - deleteOrder.verifyNoMoreInteractions(); } @Test @@ -458,7 +417,7 @@ class WorkflowServiceInstanceServiceTest { } @Test - void deleteServiceInstanceWithNoAcceptsDoesNothing() throws InterruptedException { + void deleteServiceInstanceWithNoAcceptsDoesNothing() { when(serviceInstanceStateRepository.saveState(anyString(), any(OperationState.class), anyString())) .thenReturn(Mono.just(new ServiceInstanceState(OperationState.IN_PROGRESS, "delete service instance started", new Timestamp(Instant.now().minusSeconds(60).toEpochMilli())))) @@ -475,8 +434,6 @@ class WorkflowServiceInstanceServiceTest { given(deleteServiceInstanceWorkflow2.accept(request)) .willReturn(Mono.just(false)); - InOrder deleteOrder = inOrder(deleteServiceInstanceWorkflow1, deleteServiceInstanceWorkflow2); - StepVerifier.create(workflowServiceInstanceService.deleteServiceInstance(request)) .assertNext(response -> { InOrder repoOrder = inOrder(serviceInstanceStateRepository); @@ -486,22 +443,15 @@ class WorkflowServiceInstanceServiceTest { .saveState(eq("foo"), eq(OperationState.SUCCEEDED), eq("delete service instance completed")); repoOrder.verifyNoMoreInteractions(); - deleteOrder.verify(deleteServiceInstanceWorkflow2).accept(eq(request)); - deleteOrder.verify(deleteServiceInstanceWorkflow1).accept(eq(request)); + verifyNoMoreInteractions(createServiceInstanceWorkflow1, createServiceInstanceWorkflow2); assertThat(response).isNotNull(); }) .verifyComplete(); - - Thread.sleep(50); - - deleteOrder.verify(deleteServiceInstanceWorkflow2).accept(eq(request)); - deleteOrder.verify(deleteServiceInstanceWorkflow1).accept(eq(request)); - deleteOrder.verifyNoMoreInteractions(); } @Test - void updateServiceInstance() throws InterruptedException { + void updateServiceInstance() { when(serviceInstanceStateRepository.saveState(anyString(), any(OperationState.class), anyString())) .thenReturn(Mono.just(new ServiceInstanceState(OperationState.IN_PROGRESS, "update service instance started", new Timestamp(Instant.now().minusSeconds(60).toEpochMilli())))) @@ -538,16 +488,7 @@ class WorkflowServiceInstanceServiceTest { .dashboardUrl("https://dashboard.example.com") .operation("working2"))); given(updateServiceInstanceWorkflow2.update(eq(request), eq(builtResponse))) - .willReturn(Mono.just(true) - .flatMap(value -> { - try { - Thread.sleep(500); - } catch (InterruptedException e) { - } - return higherOrderFlow.mono(); - })); - - InOrder updateOrder = inOrder(updateServiceInstanceWorkflow1, updateServiceInstanceWorkflow2); + .willReturn(higherOrderFlow.mono()); StepVerifier.create(workflowServiceInstanceService.updateServiceInstance(request)) .assertNext(response -> { @@ -558,34 +499,31 @@ class WorkflowServiceInstanceServiceTest { .saveState(eq("foo"), eq(OperationState.SUCCEEDED), eq("update service instance completed")); repoOrder.verifyNoMoreInteractions(); + lowerOrderFlow.complete(); + lowerOrderFlow.assertWasNotRequested(); + + higherOrderFlow.complete(); + lowerOrderFlow.assertWasRequested(); + + InOrder updateOrder = inOrder(updateServiceInstanceWorkflow1, updateServiceInstanceWorkflow2); updateOrder.verify(updateServiceInstanceWorkflow2).buildResponse(eq(request), any(UpdateServiceInstanceResponseBuilder.class)); updateOrder.verify(updateServiceInstanceWorkflow1).buildResponse(eq(request), any(UpdateServiceInstanceResponseBuilder.class)); + updateOrder.verify(updateServiceInstanceWorkflow2).update(request, builtResponse); + updateOrder.verify(updateServiceInstanceWorkflow1).update(request, builtResponse); + updateOrder.verifyNoMoreInteractions(); assertThat(response).isNotNull(); assertThat(response.isAsync()).isTrue(); assertThat(response.getDashboardUrl()).isEqualTo("https://dashboard.example.com"); assertThat(response.getOperation()).isEqualTo("working2"); }) - .expectComplete() - .verifyThenAssertThat() - .tookLessThan(Duration.ofMillis(250)); - - lowerOrderFlow.complete(); - lowerOrderFlow.assertWasNotRequested(); - - higherOrderFlow.complete(); - Thread.sleep(600); - lowerOrderFlow.assertWasRequested(); - - updateOrder.verify(updateServiceInstanceWorkflow2).update(request, builtResponse); - updateOrder.verify(updateServiceInstanceWorkflow1).update(request, builtResponse); - updateOrder.verifyNoMoreInteractions(); + .verifyComplete(); } @Test - void updateServiceInstanceWithAsyncError() throws InterruptedException { + void updateServiceInstanceWithAsyncError() { when(serviceInstanceStateRepository.saveState(anyString(), any(OperationState.class), anyString())) .thenReturn(Mono.just(new ServiceInstanceState(OperationState.IN_PROGRESS, "update service instance started", new Timestamp(Instant.now().minusSeconds(60).toEpochMilli())))) @@ -612,31 +550,27 @@ class WorkflowServiceInstanceServiceTest { given(updateServiceInstanceWorkflow2.buildResponse(eq(request), any(UpdateServiceInstanceResponseBuilder.class))) .willReturn(Mono.just(responseBuilder)); - InOrder repoOrder = inOrder(serviceInstanceStateRepository); - InOrder updateOrder = inOrder(updateServiceInstanceWorkflow1, updateServiceInstanceWorkflow2); - StepVerifier.create(workflowServiceInstanceService.updateServiceInstance(request)) .assertNext(response -> { + InOrder repoOrder = inOrder(serviceInstanceStateRepository); repoOrder.verify(serviceInstanceStateRepository) .saveState(eq("foo"), eq(OperationState.IN_PROGRESS), eq("update service instance started")); + repoOrder.verify(serviceInstanceStateRepository) + .saveState(eq("foo"), eq(OperationState.FAILED), eq("update foo error")); + repoOrder.verifyNoMoreInteractions(); + InOrder updateOrder = inOrder(updateServiceInstanceWorkflow1, updateServiceInstanceWorkflow2); updateOrder.verify(updateServiceInstanceWorkflow2).buildResponse(eq(request), any(UpdateServiceInstanceResponseBuilder.class)); updateOrder.verify(updateServiceInstanceWorkflow1).buildResponse(eq(request), any(UpdateServiceInstanceResponseBuilder.class)); + updateOrder.verify(updateServiceInstanceWorkflow2).update(request, responseBuilder.build()); + updateOrder.verify(updateServiceInstanceWorkflow1).update(request, responseBuilder.build()); + updateOrder.verifyNoMoreInteractions(); assertThat(response).isNotNull(); }) .verifyComplete(); - - Thread.sleep(50); - - repoOrder.verify(serviceInstanceStateRepository).saveState(eq("foo"), eq(OperationState.FAILED), eq("update foo error")); - repoOrder.verifyNoMoreInteractions(); - - updateOrder.verify(updateServiceInstanceWorkflow2).update(request, responseBuilder.build()); - updateOrder.verify(updateServiceInstanceWorkflow1).update(request, responseBuilder.build()); - updateOrder.verifyNoMoreInteractions(); } @Test @@ -665,7 +599,7 @@ class WorkflowServiceInstanceServiceTest { } @Test - void updateServiceInstanceWithNoAcceptsDoesNothing() throws InterruptedException { + void updateServiceInstanceWithNoAcceptsDoesNothing() { when(serviceInstanceStateRepository.saveState(anyString(), any(OperationState.class), anyString())) .thenReturn(Mono.just(new ServiceInstanceState(OperationState.IN_PROGRESS, "update service instance started", new Timestamp(Instant.now().minusSeconds(60).toEpochMilli())))) @@ -682,8 +616,6 @@ class WorkflowServiceInstanceServiceTest { given(updateServiceInstanceWorkflow2.accept(request)) .willReturn(Mono.just(false)); - InOrder updateOrder = inOrder(updateServiceInstanceWorkflow1, updateServiceInstanceWorkflow2); - StepVerifier.create(workflowServiceInstanceService.updateServiceInstance(request)) .assertNext(response -> { InOrder repoOrder = inOrder(serviceInstanceStateRepository); @@ -693,18 +625,11 @@ class WorkflowServiceInstanceServiceTest { .saveState(eq("foo"), eq(OperationState.SUCCEEDED), eq("update service instance completed")); repoOrder.verifyNoMoreInteractions(); - updateOrder.verify(updateServiceInstanceWorkflow2).accept(eq(request)); - updateOrder.verify(updateServiceInstanceWorkflow1).accept(eq(request)); + verifyNoMoreInteractions(updateServiceInstanceWorkflow1, updateServiceInstanceWorkflow2); assertThat(response).isNotNull(); }) .verifyComplete(); - - Thread.sleep(50); - - updateOrder.verify(updateServiceInstanceWorkflow2).accept(eq(request)); - updateOrder.verify(updateServiceInstanceWorkflow1).accept(eq(request)); - updateOrder.verifyNoMoreInteractions(); } @Order(Ordered.HIGHEST_PRECEDENCE)