From 86c510a8c59645b618e92b432d17b69e634b5bf0 Mon Sep 17 00:00:00 2001 From: Mark Paluch Date: Thu, 1 Jun 2017 10:04:33 +0200 Subject: [PATCH] DATACASS-453 - Upgrade to Reactor 3.1 M2. Adopt to API change from Publisher.subscribe() to Publisher.toProcessor(). Adopt to changed reactor-test groupId. Provide mocks for calls that allowed previously null Publishers. --- spring-data-cassandra/pom.xml | 2 +- .../query/ReactiveCassandraParameterAccessor.java | 4 ++-- .../core/DefaultBridgedReactiveSessionUnitTests.java | 12 +++++++++++- 3 files changed, 14 insertions(+), 4 deletions(-) diff --git a/spring-data-cassandra/pom.xml b/spring-data-cassandra/pom.xml index 1eada1568..9f049b494 100644 --- a/spring-data-cassandra/pom.xml +++ b/spring-data-cassandra/pom.xml @@ -90,7 +90,7 @@ - io.projectreactor.addons + io.projectreactor reactor-test ${reactor} test diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/query/ReactiveCassandraParameterAccessor.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/query/ReactiveCassandraParameterAccessor.java index 536ebf41d..fe7dad3a0 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/query/ReactiveCassandraParameterAccessor.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/query/ReactiveCassandraParameterAccessor.java @@ -52,9 +52,9 @@ class ReactiveCassandraParameterAccessor extends CassandraParametersParameterAcc } if (ReactiveWrappers.isSingleValueType(value.getClass())) { - subscriptions.add(ReactiveWrapperConverters.toWrapper(value, Mono.class).subscribe()); + subscriptions.add(ReactiveWrapperConverters.toWrapper(value, Mono.class).toProcessor()); } else { - subscriptions.add(ReactiveWrapperConverters.toWrapper(value, Flux.class).collectList().subscribe()); + subscriptions.add(ReactiveWrapperConverters.toWrapper(value, Flux.class).collectList().toProcessor()); } } } diff --git a/spring-data-cassandra/src/test/java/org/springframework/data/cql/core/DefaultBridgedReactiveSessionUnitTests.java b/spring-data-cassandra/src/test/java/org/springframework/data/cql/core/DefaultBridgedReactiveSessionUnitTests.java index 41a03ca5b..2973ae8cd 100644 --- a/spring-data-cassandra/src/test/java/org/springframework/data/cql/core/DefaultBridgedReactiveSessionUnitTests.java +++ b/spring-data-cassandra/src/test/java/org/springframework/data/cql/core/DefaultBridgedReactiveSessionUnitTests.java @@ -31,9 +31,13 @@ import org.mockito.junit.MockitoJUnitRunner; import org.springframework.data.cql.core.session.DefaultBridgedReactiveSession; import com.datastax.driver.core.Cluster; +import com.datastax.driver.core.PreparedStatement; +import com.datastax.driver.core.RegularStatement; +import com.datastax.driver.core.ResultSetFuture; import com.datastax.driver.core.Session; import com.datastax.driver.core.SimpleStatement; import com.datastax.driver.core.Statement; +import com.google.common.util.concurrent.ListenableFuture; /** * Unit tests for {@link DefaultBridgedReactiveSession}. @@ -43,13 +47,19 @@ import com.datastax.driver.core.Statement; @RunWith(MockitoJUnitRunner.class) public class DefaultBridgedReactiveSessionUnitTests { - @Mock private Session sessionMock; + @Mock Session sessionMock; + @Mock ResultSetFuture future; + @Mock ListenableFuture preparedStatementFuture; private DefaultBridgedReactiveSession reactiveSession; @Before public void before() throws Exception { + reactiveSession = new DefaultBridgedReactiveSession(sessionMock, Schedulers.immediate()); + + when(sessionMock.executeAsync(any(Statement.class))).thenReturn(future); + when(sessionMock.prepareAsync(any(RegularStatement.class))).thenReturn(preparedStatementFuture); } @Test // DATACASS-335