Revert "Fix blocking calls to workflow layers and return a response as soon as possible"
This reverts commit 9072f85836.
This commit is contained in:
committed by
Alberto Ríos
parent
e9d3dd1b8c
commit
fd245587af
@@ -80,6 +80,7 @@ public class WorkflowServiceInstanceBindingService implements ServiceInstanceBin
|
||||
@Override
|
||||
public Mono<CreateServiceInstanceBindingResponse> 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<DeleteServiceInstanceBindingResponse> 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"))
|
||||
|
||||
@@ -77,6 +77,7 @@ public class WorkflowServiceInstanceService implements ServiceInstanceService {
|
||||
@Override
|
||||
public Mono<CreateServiceInstanceResponse> 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<DeleteServiceInstanceResponse> 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<Void> 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<UpdateServiceInstanceResponse> 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<Void> 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"))
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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)
|
||||
|
||||
Reference in New Issue
Block a user