DATACASS-368 - Migrate reactive tests from TestSubscriber to StepVerifier.

Replace TestSubscriber and .block() calls in test with StepVerifier.
This commit is contained in:
Mark Paluch
2017-03-24 14:38:56 +01:00
parent ed7565273e
commit acbb3ba5f9
11 changed files with 353 additions and 1628 deletions

View File

@@ -53,6 +53,13 @@
<optional>true</optional>
</dependency>
<dependency>
<groupId>io.projectreactor.addons</groupId>
<artifactId>reactor-test</artifactId>
<version>${reactor}</version>
<scope>test</scope>
</dependency>
<dependency>
<groupId>com.datastax.cassandra</groupId>
<artifactId>cassandra-driver-core</artifactId>

View File

@@ -19,6 +19,7 @@ import static org.assertj.core.api.Assertions.*;
import reactor.core.publisher.Mono;
import reactor.core.scheduler.Schedulers;
import reactor.test.StepVerifier;
import org.junit.Before;
import org.junit.Test;
@@ -27,8 +28,6 @@ import org.springframework.cassandra.core.session.ReactiveResultSet;
import org.springframework.cassandra.test.integration.AbstractKeyspaceCreatingIntegrationTest;
import com.datastax.driver.core.KeyspaceMetadata;
import com.datastax.driver.core.PreparedStatement;
import com.datastax.driver.core.Row;
import com.datastax.driver.core.exceptions.SyntaxError;
/**
@@ -38,7 +37,7 @@ import com.datastax.driver.core.exceptions.SyntaxError;
*/
public class DefaultBridgedReactiveSessionIntegrationTests extends AbstractKeyspaceCreatingIntegrationTest {
private DefaultBridgedReactiveSession reactiveSession;
DefaultBridgedReactiveSession reactiveSession;
@Before
public void before() throws Exception {
@@ -59,18 +58,17 @@ public class DefaultBridgedReactiveSessionIntegrationTests extends AbstractKeysp
assertThat(keyspace.getTable("users")).isNull();
ReactiveResultSet resultSet = execution.block();
StepVerifier.create(execution).consumeNextWith(actual -> {
assertThat(actual.wasApplied()).isTrue();
}).verifyComplete();
assertThat(resultSet.wasApplied()).isTrue();
assertThat(keyspace.getTable("users")).isNotNull();
}
@Test(expected = SyntaxError.class) // DATACASS-335
public void executeShouldTransportExceptionsInMono() throws Exception {
Mono<ReactiveResultSet> execution = reactiveSession.execute("INSERT INTO dummy;");
execution.block();
@Test // DATACASS-335
public void executeShouldTransportExceptionsInMono() {
StepVerifier.create(reactiveSession.execute("INSERT INTO dummy;")).expectError(SyntaxError.class).verify();
}
@Test // DATACASS-335
@@ -79,12 +77,13 @@ public class DefaultBridgedReactiveSessionIntegrationTests extends AbstractKeysp
session.execute("CREATE TABLE users (\n" + " userid text PRIMARY KEY,\n" + " first_name text\n" + ");");
session.execute("INSERT INTO users (userid, first_name) VALUES ('White', 'Walter');");
Mono<ReactiveResultSet> execution = reactiveSession.execute("SELECT * FROM users;");
ReactiveResultSet resultSet = execution.block();
Row row = resultSet.rows().blockFirst();
StepVerifier.create(reactiveSession.execute("SELECT * FROM users;")).consumeNextWith(actual -> {
assertThat(row).isNotNull();
assertThat(row.getString("userid")).isEqualTo("White");
StepVerifier.create(actual.rows()).consumeNextWith(row -> {
assertThat(row.getString("userid")).isEqualTo("White");
}).verifyComplete();
}).verifyComplete();
}
@Test // DATACASS-335
@@ -92,12 +91,11 @@ public class DefaultBridgedReactiveSessionIntegrationTests extends AbstractKeysp
session.execute("CREATE TABLE users (\n" + " userid text PRIMARY KEY,\n" + " first_name text\n" + ");");
Mono<PreparedStatement> execution = reactiveSession
.prepare("INSERT INTO users (userid, first_name) VALUES (?, ?);");
PreparedStatement preparedStatement = execution.block();
StepVerifier.create(reactiveSession.prepare("INSERT INTO users (userid, first_name) VALUES (?, ?);"))
.consumeNextWith(actual -> {
assertThat(preparedStatement).isNotNull();
assertThat(preparedStatement.getQueryString()).isEqualTo("INSERT INTO users (userid, first_name) VALUES (?, ?);");
assertThat(actual.getQueryString()).isEqualTo("INSERT INTO users (userid, first_name) VALUES (?, ?);");
}).verifyComplete();
}
private KeyspaceMetadata getKeyspaceMetadata() {

View File

@@ -18,8 +18,8 @@ package org.springframework.cassandra.core;
import static org.assertj.core.api.Assertions.*;
import reactor.core.scheduler.Schedulers;
import reactor.test.StepVerifier;
import java.util.Map;
import java.util.concurrent.atomic.AtomicBoolean;
import org.junit.Before;
@@ -39,11 +39,12 @@ import com.datastax.driver.core.querybuilder.QueryBuilder;
public class ReactiveCqlTemplateIntegrationTests extends AbstractKeyspaceCreatingIntegrationTest {
private static final AtomicBoolean initialized = new AtomicBoolean();
private ReactiveSession reactiveSession;
private ReactiveCqlTemplate template;
ReactiveSession reactiveSession;
ReactiveCqlTemplate template;
@Before
public void before() throws Exception {
public void before() {
reactiveSession = new DefaultBridgedReactiveSession(getSession(), Schedulers.elastic());
@@ -59,74 +60,88 @@ public class ReactiveCqlTemplateIntegrationTests extends AbstractKeyspaceCreatin
}
@Test // DATACASS-335
public void executeShouldRemoveRecords() throws Exception {
public void executeShouldRemoveRecords() {
template.execute("DELETE FROM user WHERE id = 'WHITE'").block();
StepVerifier.create(template.execute("DELETE FROM user WHERE id = 'WHITE'")).expectNext(true).verifyComplete();
assertThat(getSession().execute("SELECT * FROM user").one()).isNull();
}
@Test // DATACASS-335
public void queryForObjectShouldReturnFirstColumn() throws Exception {
public void queryForObjectShouldReturnFirstColumn() {
String id = template.queryForObject("SELECT id FROM user;", String.class).block();
assertThat(id).isEqualTo("WHITE");
StepVerifier.create(template.queryForObject("SELECT id FROM user;", String.class)) //
.expectNext("WHITE") //
.verifyComplete();
}
@Test // DATACASS-335
public void queryForObjectShouldReturnMap() throws Exception {
public void queryForObjectShouldReturnMap() {
Map<String, Object> map = template.queryForMap("SELECT * FROM user;").block();
StepVerifier.create(template.queryForMap("SELECT * FROM user;")) //
.consumeNextWith(actual -> {
assertThat(map).containsEntry("id", "WHITE").containsEntry("username", "Walter");
assertThat(actual).containsEntry("id", "WHITE").containsEntry("username", "Walter");
}).verifyComplete();
}
@Test // DATACASS-335
public void executeStatementShouldRemoveRecords() throws Exception {
public void executeStatementShouldRemoveRecords() {
template.execute(QueryBuilder.delete().from("user").where(QueryBuilder.eq("id", "WHITE"))).block();
StepVerifier
.create(template.execute(QueryBuilder.delete() //
.from("user") //
.where(QueryBuilder.eq("id", "WHITE")))) //
.expectNext(true) //
.verifyComplete();
assertThat(getSession().execute("SELECT * FROM user").one()).isNull();
}
@Test // DATACASS-335
public void queryForObjectStatementShouldReturnFirstColumn() throws Exception {
public void queryForObjectStatementShouldReturnFirstColumn() {
String id = template.queryForObject(QueryBuilder.select("id").from("user"), String.class).block();
assertThat(id).isEqualTo("WHITE");
StepVerifier
.create(template.queryForObject(QueryBuilder //
.select("id") //
.from("user"), String.class)) //
.expectNext("WHITE") //
.verifyComplete();
}
@Test // DATACASS-335
public void queryForObjectStatementShouldReturnMap() throws Exception {
public void queryForObjectStatementShouldReturnMap() {
Map<String, Object> map = template.queryForMap(QueryBuilder.select().from("user")).block();
StepVerifier.create(template.queryForMap(QueryBuilder.select().from("user"))) //
.consumeNextWith(actual -> {
assertThat(map).containsEntry("id", "WHITE").containsEntry("username", "Walter");
assertThat(actual).containsEntry("id", "WHITE").containsEntry("username", "Walter");
}).verifyComplete();
}
@Test // DATACASS-335
public void executeWithArgsShouldRemoveRecords() throws Exception {
public void executeWithArgsShouldRemoveRecords() {
template.execute("DELETE FROM user WHERE id = ?", "WHITE").block();
StepVerifier.create(template.execute("DELETE FROM user WHERE id = ?", "WHITE")).expectNext(true).verifyComplete();
assertThat(getSession().execute("SELECT * FROM user").one()).isNull();
}
@Test // DATACASS-335
public void queryForObjectWithArgsShouldReturnFirstColumn() throws Exception {
public void queryForObjectWithArgsShouldReturnFirstColumn() {
String id = template.queryForObject("SELECT id FROM user WHERE id = ?;", String.class, "WHITE").block();
assertThat(id).isEqualTo("WHITE");
StepVerifier.create(template.queryForObject("SELECT id FROM user WHERE id = ?;", String.class, "WHITE")) //
.expectNext("WHITE") //
.verifyComplete();
}
@Test // DATACASS-335
public void queryForObjectWithArgsShouldReturnMap() throws Exception {
public void queryForObjectWithArgsShouldReturnMap() {
Map<String, Object> map = template.queryForMap("SELECT * FROM user WHERE id = ?;", "WHITE").block();
StepVerifier.create(template.queryForMap("SELECT * FROM user WHERE id = ?;", "WHITE")) //
.consumeNextWith(actual -> {
assertThat(map).containsEntry("id", "WHITE").containsEntry("username", "Walter");
assertThat(actual).containsEntry("id", "WHITE").containsEntry("username", "Walter");
}).verifyComplete();
}
}

View File

@@ -22,9 +22,9 @@ import static org.mockito.Mockito.*;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import reactor.test.StepVerifier;
import java.util.Collections;
import java.util.List;
import java.util.function.Consumer;
import org.junit.Before;
@@ -90,7 +90,8 @@ public class ReactiveCqlTemplateUnitTests {
});
verify(session, never()).close();
assertThat(flux.blockLast()).isEqualTo("OK");
StepVerifier.create(flux).expectNext("OK").verifyComplete();
verify(session).close();
}
@@ -101,13 +102,7 @@ public class ReactiveCqlTemplateUnitTests {
throw new InvalidQueryException("wrong query");
});
try {
flux.blockLast();
fail("Missing CassandraInvalidQueryException");
} catch (CassandraInvalidQueryException e) {
assertThat(e).hasMessageContaining("wrong query");
}
StepVerifier.create(flux).expectError(CassandraInvalidQueryException.class).verify();
}
@Test // DATACASS-335
@@ -118,7 +113,9 @@ public class ReactiveCqlTemplateUnitTests {
Mono<Boolean> mono = template.execute("UPDATE user SET a = 'b';");
verifyZeroInteractions(session);
assertThat(mono.block()).isFalse();
StepVerifier.create(mono).expectNext(false).verifyComplete();
verify(session).execute(any(Statement.class));
}
@@ -129,13 +126,7 @@ public class ReactiveCqlTemplateUnitTests {
Mono<Boolean> mono = template.execute("UPDATE user SET a = 'b';");
try {
mono.block();
fail("Missing CassandraConnectionFailureException");
} catch (CassandraConnectionFailureException e) {
assertThat(e).hasMessageContaining("tried for query failed");
}
StepVerifier.create(mono).expectError(CassandraConnectionFailureException.class).verify();
}
// -------------------------------------------------------------------------
@@ -147,7 +138,7 @@ public class ReactiveCqlTemplateUnitTests {
doTestStrings(null, null, null, reactiveCqlTemplate -> {
reactiveCqlTemplate.execute("SELECT * from USERS").block();
StepVerifier.create(reactiveCqlTemplate.execute("SELECT * from USERS")).expectNextCount(1).verifyComplete();
verify(session).execute(any(Statement.class));
});
@@ -158,7 +149,9 @@ public class ReactiveCqlTemplateUnitTests {
doTestStrings(5, ConsistencyLevel.ONE, DowngradingConsistencyRetryPolicy.INSTANCE, reactiveCqlTemplate -> {
reactiveCqlTemplate.execute("SELECT * from USERS").block();
StepVerifier.create(reactiveCqlTemplate.execute("SELECT * from USERS")) //
.expectNextCount(1) //
.verifyComplete();
verify(session).execute(any(Statement.class));
});
@@ -171,9 +164,8 @@ public class ReactiveCqlTemplateUnitTests {
Mono<ReactiveResultSet> mono = reactiveCqlTemplate.queryForResultSet("SELECT * from USERS");
List<Row> rows = mono.block().rows().collectList().block();
StepVerifier.create(mono.flatMap(ReactiveResultSet::rows)).expectNextCount(3).verifyComplete();
assertThat(rows).hasSize(3);
verify(session).execute(any(Statement.class));
});
}
@@ -185,9 +177,8 @@ public class ReactiveCqlTemplateUnitTests {
Flux<String> flux = reactiveCqlTemplate.query("SELECT * from USERS", (row, index) -> row.getString(0));
List<String> rows = flux.collectList().block();
StepVerifier.create(flux).expectNext("Walter", "Hank", " Jesse").verifyComplete();
assertThat(rows).hasSize(3).contains("Walter", "Hank", " Jesse");
verify(session).execute(any(Statement.class));
});
}
@@ -199,9 +190,8 @@ public class ReactiveCqlTemplateUnitTests {
Flux<String> flux = reactiveCqlTemplate.query("SELECT * from USERS", (row, index) -> row.getString(0));
List<String> rows = flux.collectList().block();
StepVerifier.create(flux).expectNext("Walter", "Hank", " Jesse").verifyComplete();
assertThat(rows).hasSize(3).contains("Walter", "Hank", " Jesse");
verify(session).execute(any(Statement.class));
});
}
@@ -215,7 +205,9 @@ public class ReactiveCqlTemplateUnitTests {
Flux<Boolean> flux = template.query("UPDATE user SET a = 'b';", resultSet -> Mono.just(resultSet.wasApplied()));
verifyZeroInteractions(session);
assertThat(flux.collectList().block()).hasSize(1).contains(true);
StepVerifier.create(flux).expectNext(true).verifyComplete();
verify(session).execute(any(Statement.class));
}
@@ -226,13 +218,7 @@ public class ReactiveCqlTemplateUnitTests {
Flux<Boolean> flux = template.query("UPDATE user SET a = 'b';", resultSet -> Mono.just(resultSet.wasApplied()));
try {
flux.blockLast();
fail("Missing CassandraConnectionFailureException");
} catch (CassandraConnectionFailureException e) {
assertThat(e).hasMessageContaining("tried for query failed");
}
StepVerifier.create(flux).expectError(CassandraConnectionFailureException.class).verify();
}
@Test // DATACASS-335
@@ -242,7 +228,8 @@ public class ReactiveCqlTemplateUnitTests {
when(reactiveResultSet.rows()).thenReturn(Flux.empty());
Mono<String> mono = template.queryForObject("SELECT * FROM user", (row, rowNum) -> "OK");
assertThat(mono.hasElement().block()).isFalse();
StepVerifier.create(mono).verifyComplete();
}
@Test // DATACASS-335
@@ -252,7 +239,8 @@ public class ReactiveCqlTemplateUnitTests {
when(reactiveResultSet.rows()).thenReturn(Flux.just(row));
Mono<String> mono = template.queryForObject("SELECT * FROM user", (row, rowNum) -> "OK");
assertThat(mono.block()).isEqualTo("OK");
StepVerifier.create(mono).expectNext("OK").verifyComplete();
}
@Test // DATACASS-335
@@ -262,7 +250,8 @@ public class ReactiveCqlTemplateUnitTests {
when(reactiveResultSet.rows()).thenReturn(Flux.just(row));
Mono<String> mono = template.queryForObject("SELECT * FROM user", (row, rowNum) -> null);
assertThat(mono.hasElement().block()).isFalse();
StepVerifier.create(mono).verifyComplete();
}
@Test // DATACASS-335
@@ -273,13 +262,7 @@ public class ReactiveCqlTemplateUnitTests {
Mono<String> mono = template.queryForObject("SELECT * FROM user", (row, rowNum) -> "OK");
try {
mono.block();
fail("Missing IncorrectResultSizeDataAccessException");
} catch (IncorrectResultSizeDataAccessException e) {
assertThat(e).hasMessageContaining("expected 1, actual 2");
}
StepVerifier.create(mono).expectError(IncorrectResultSizeDataAccessException.class).verify();
}
@Test // DATACASS-335
@@ -293,7 +276,7 @@ public class ReactiveCqlTemplateUnitTests {
Mono<String> mono = template.queryForObject("SELECT * FROM user", String.class);
assertThat(mono.block()).isEqualTo("OK");
StepVerifier.create(mono).expectNext("OK").verifyComplete();
}
@Test // DATACASS-335
@@ -307,7 +290,7 @@ public class ReactiveCqlTemplateUnitTests {
Flux<String> flux = template.queryForFlux("SELECT * FROM user", String.class);
assertThat(flux.collectList().block()).contains("OK", "NOT OK");
StepVerifier.create(flux).expectNext("OK", "NOT OK").verifyComplete();
}
@Test // DATACASS-335
@@ -318,7 +301,7 @@ public class ReactiveCqlTemplateUnitTests {
Flux<Row> flux = template.queryForRows("SELECT * FROM user");
assertThat(flux.collectList().block()).hasSize(2).contains(row);
StepVerifier.create(flux).expectNext(row, row).verifyComplete();
}
@Test // DATACASS-335
@@ -329,7 +312,7 @@ public class ReactiveCqlTemplateUnitTests {
Mono<Boolean> mono = template.execute("UPDATE user SET a = 'b';");
assertThat(mono.block()).isTrue();
StepVerifier.create(mono).expectNext(true).verifyComplete();
}
@Test // DATACASS-335
@@ -341,7 +324,9 @@ public class ReactiveCqlTemplateUnitTests {
Flux<Boolean> flux = template.execute(Flux.just("UPDATE user SET a = 'b';", "UPDATE user SET x = 'y';"));
verifyZeroInteractions(session);
assertThat(flux.collectList().block()).hasSize(2).contains(true, false);
StepVerifier.create(flux).expectNext(true).expectNext(false).verifyComplete();
verify(session, times(2)).execute(any(Statement.class));
}
@@ -354,7 +339,9 @@ public class ReactiveCqlTemplateUnitTests {
doTestStrings(null, null, null, reactiveCqlTemplate -> {
reactiveCqlTemplate.execute(new SimpleStatement("SELECT * from USERS")).block();
StepVerifier.create(reactiveCqlTemplate.execute(new SimpleStatement("SELECT * from USERS"))) //
.expectNextCount(1) //
.verifyComplete();
verify(session).execute(any(Statement.class));
});
@@ -365,7 +352,9 @@ public class ReactiveCqlTemplateUnitTests {
doTestStrings(5, ConsistencyLevel.ONE, DowngradingConsistencyRetryPolicy.INSTANCE, reactiveCqlTemplate -> {
reactiveCqlTemplate.execute(new SimpleStatement("SELECT * from USERS")).block();
StepVerifier.create(reactiveCqlTemplate.execute(new SimpleStatement("SELECT * from USERS"))) //
.expectNextCount(1) //
.verifyComplete();
verify(session).execute(any(Statement.class));
});
@@ -376,11 +365,12 @@ public class ReactiveCqlTemplateUnitTests {
doTestStrings(null, null, null, reactiveCqlTemplate -> {
Mono<ReactiveResultSet> mono = reactiveCqlTemplate.queryForResultSet(new SimpleStatement("SELECT * from USERS"));
StepVerifier
.create(reactiveCqlTemplate.queryForResultSet(new SimpleStatement("SELECT * from USERS"))
.flatMap(ReactiveResultSet::rows)) //
.expectNextCount(3) //
.verifyComplete();
List<Row> rows = mono.block().rows().collectList().block();
assertThat(rows).hasSize(3);
verify(session).execute(any(Statement.class));
});
}
@@ -393,9 +383,8 @@ public class ReactiveCqlTemplateUnitTests {
Flux<String> flux = reactiveCqlTemplate.query(new SimpleStatement("SELECT * from USERS"),
(row, index) -> row.getString(0));
List<String> rows = flux.collectList().block();
StepVerifier.create(flux).expectNext("Walter", "Hank", " Jesse").verifyComplete();
assertThat(rows).hasSize(3).contains("Walter", "Hank", " Jesse");
verify(session).execute(any(Statement.class));
});
}
@@ -408,9 +397,11 @@ public class ReactiveCqlTemplateUnitTests {
Flux<String> flux = reactiveCqlTemplate.query(new SimpleStatement("SELECT * from USERS"),
(row, index) -> row.getString(0));
List<String> rows = flux.collectList().block();
StepVerifier.create(flux.collectList()).consumeNextWith(rows -> {
assertThat(rows).hasSize(3).contains("Walter", "Hank", " Jesse");
}).verifyComplete();
assertThat(rows).hasSize(3).contains("Walter", "Hank", " Jesse");
verify(session).execute(any(Statement.class));
});
}
@@ -425,7 +416,7 @@ public class ReactiveCqlTemplateUnitTests {
resultSet -> Mono.just(resultSet.wasApplied()));
verifyZeroInteractions(session);
assertThat(flux.collectList().block()).hasSize(1).contains(true);
StepVerifier.create(flux).expectNext(true).verifyComplete();
verify(session).execute(any(Statement.class));
}
@@ -437,13 +428,7 @@ public class ReactiveCqlTemplateUnitTests {
Flux<Boolean> flux = template.query(new SimpleStatement("UPDATE user SET a = 'b';"),
resultSet -> Mono.just(resultSet.wasApplied()));
try {
flux.blockLast();
fail("Missing CassandraConnectionFailureException");
} catch (CassandraConnectionFailureException e) {
assertThat(e).hasMessageContaining("tried for query failed");
}
StepVerifier.create(flux).expectError(CassandraConnectionFailureException.class).verify();
}
@Test // DATACASS-335
@@ -453,7 +438,8 @@ public class ReactiveCqlTemplateUnitTests {
when(reactiveResultSet.rows()).thenReturn(Flux.empty());
Mono<String> mono = template.queryForObject(new SimpleStatement("SELECT * FROM user"), (row, rowNum) -> "OK");
assertThat(mono.hasElement().block()).isFalse();
StepVerifier.create(mono).verifyComplete();
}
@Test // DATACASS-335
@@ -463,7 +449,8 @@ public class ReactiveCqlTemplateUnitTests {
when(reactiveResultSet.rows()).thenReturn(Flux.just(row));
Mono<String> mono = template.queryForObject(new SimpleStatement("SELECT * FROM user"), (row, rowNum) -> "OK");
assertThat(mono.block()).isEqualTo("OK");
StepVerifier.create(mono).expectNext("OK").verifyComplete();
}
@Test // DATACASS-335
@@ -473,7 +460,8 @@ public class ReactiveCqlTemplateUnitTests {
when(reactiveResultSet.rows()).thenReturn(Flux.just(row));
Mono<String> mono = template.queryForObject(new SimpleStatement("SELECT * FROM user"), (row, rowNum) -> null);
assertThat(mono.hasElement().block()).isFalse();
StepVerifier.create(mono).verifyComplete();
}
@Test // DATACASS-335
@@ -484,13 +472,7 @@ public class ReactiveCqlTemplateUnitTests {
Mono<String> mono = template.queryForObject(new SimpleStatement("SELECT * FROM user"), (row, rowNum) -> "OK");
try {
mono.block();
fail("Missing IncorrectResultSizeDataAccessException");
} catch (IncorrectResultSizeDataAccessException e) {
assertThat(e).hasMessageContaining("expected 1, actual 2");
}
StepVerifier.create(mono).expectError(IncorrectResultSizeDataAccessException.class).verify();
}
@Test // DATACASS-335
@@ -504,7 +486,7 @@ public class ReactiveCqlTemplateUnitTests {
Mono<String> mono = template.queryForObject(new SimpleStatement("SELECT * FROM user"), String.class);
assertThat(mono.block()).isEqualTo("OK");
StepVerifier.create(mono).expectNext("OK").verifyComplete();
}
@Test // DATACASS-335
@@ -518,7 +500,7 @@ public class ReactiveCqlTemplateUnitTests {
Flux<String> flux = template.queryForFlux(new SimpleStatement("SELECT * FROM user"), String.class);
assertThat(flux.collectList().block()).contains("OK", "NOT OK");
StepVerifier.create(flux).expectNext("OK", "NOT OK").verifyComplete();
}
@Test // DATACASS-335
@@ -529,7 +511,7 @@ public class ReactiveCqlTemplateUnitTests {
Flux<Row> flux = template.queryForRows(new SimpleStatement("SELECT * FROM user"));
assertThat(flux.collectList().block()).hasSize(2).contains(row);
StepVerifier.create(flux).expectNext(row, row).verifyComplete();
}
@Test // DATACASS-335
@@ -538,9 +520,8 @@ public class ReactiveCqlTemplateUnitTests {
when(session.execute(any(Statement.class))).thenReturn(Mono.just(reactiveResultSet));
when(reactiveResultSet.wasApplied()).thenReturn(true);
Mono<Boolean> mono = template.execute(new SimpleStatement("UPDATE user SET a = 'b';"));
assertThat(mono.block()).isTrue();
StepVerifier.create(template.execute(new SimpleStatement("UPDATE user SET a = 'b';"))).expectNext(true)
.verifyComplete();
}
// -------------------------------------------------------------------------
@@ -557,9 +538,7 @@ public class ReactiveCqlTemplateUnitTests {
return session.execute(ps.bind("A")).flatMap(ReactiveResultSet::rows);
});
List<Row> rows = flux.collectList().block();
assertThat(rows).hasSize(3);
StepVerifier.create(flux).expectNextCount(3).verifyComplete();
});
}
@@ -572,7 +551,7 @@ public class ReactiveCqlTemplateUnitTests {
when(this.preparedStatement.bind("White")).thenReturn(this.boundStatement);
when(this.reactiveResultSet.wasApplied()).thenReturn(true);
assertThat(applied.block()).isTrue();
StepVerifier.create(applied).expectNext(true).verifyComplete();
});
}
@@ -587,7 +566,9 @@ public class ReactiveCqlTemplateUnitTests {
(session, ps) -> session.execute(ps.bind()));
verifyZeroInteractions(session);
assertThat(flux.collectList().block()).hasSize(1).contains(reactiveResultSet);
StepVerifier.create(flux).expectNext(reactiveResultSet).verifyComplete();
verify(session).prepare(anyString());
verify(session).execute(boundStatement);
}
@@ -601,7 +582,9 @@ public class ReactiveCqlTemplateUnitTests {
(session, ps) -> session.execute(boundStatement));
verifyZeroInteractions(session);
assertThat(flux.collectList().block()).hasSize(1).contains(reactiveResultSet);
StepVerifier.create(flux).expectNext(reactiveResultSet).verifyComplete();
verify(session).execute(boundStatement);
}
@@ -612,13 +595,7 @@ public class ReactiveCqlTemplateUnitTests {
throw new NoHostAvailableException(Collections.emptyMap());
}, (session, ps) -> session.execute(boundStatement));
try {
flux.blockLast();
fail("Missing CassandraConnectionFailureException");
} catch (CassandraConnectionFailureException e) {
assertThat(e).hasMessageContaining("tried for query");
}
StepVerifier.create(flux).expectError(CassandraConnectionFailureException.class).verify();
}
@Test // DATACASS-335
@@ -628,13 +605,7 @@ public class ReactiveCqlTemplateUnitTests {
throw new NoHostAvailableException(Collections.emptyMap());
});
try {
flux.blockLast();
fail("Missing CassandraConnectionFailureException");
} catch (CassandraConnectionFailureException e) {
assertThat(e).hasMessageContaining("tried for query");
}
StepVerifier.create(flux).expectError(CassandraConnectionFailureException.class).verify();
}
@Test // DATACASS-335
@@ -647,7 +618,8 @@ public class ReactiveCqlTemplateUnitTests {
Flux<Row> flux = template.query(session -> Mono.just(preparedStatement), ReactiveResultSet::rows);
verifyZeroInteractions(session);
assertThat(flux.collectList().block()).hasSize(1).contains(row);
StepVerifier.create(flux).expectNext(row).verifyComplete();
verify(preparedStatement).bind();
}
@@ -663,7 +635,9 @@ public class ReactiveCqlTemplateUnitTests {
}, ReactiveResultSet::rows);
verifyZeroInteractions(session);
assertThat(flux.collectList().block()).hasSize(1).contains(row);
StepVerifier.create(flux).expectNext(row).verifyComplete();
verify(preparedStatement).bind("a", "b");
}
@@ -679,7 +653,9 @@ public class ReactiveCqlTemplateUnitTests {
}, (row, rowNum) -> row);
verifyZeroInteractions(session);
assertThat(flux.collectList().block()).hasSize(1).contains(row);
StepVerifier.create(flux).expectNext(row).verifyComplete();
verify(preparedStatement).bind("a", "b");
}
@@ -693,7 +669,8 @@ public class ReactiveCqlTemplateUnitTests {
Mono<String> mono = template.queryForObject("SELECT * FROM user WHERE username = ?", (row, rowNum) -> "OK",
"Walter");
assertThat(mono.hasElement().block()).isFalse();
StepVerifier.create(mono).verifyComplete();
}
@Test // DATACASS-335
@@ -706,7 +683,8 @@ public class ReactiveCqlTemplateUnitTests {
Mono<String> mono = template.queryForObject("SELECT * FROM user WHERE username = ?", (row, rowNum) -> "OK",
"Walter");
assertThat(mono.block()).isEqualTo("OK");
StepVerifier.create(mono).expectNext("OK").verifyComplete();
}
@Test // DATACASS-335
@@ -719,13 +697,8 @@ public class ReactiveCqlTemplateUnitTests {
Mono<String> mono = template.queryForObject("SELECT * FROM user WHERE username = ?", (row, rowNum) -> "OK",
"Walter");
try {
mono.block();
fail("Missing IncorrectResultSizeDataAccessException");
} catch (IncorrectResultSizeDataAccessException e) {
assertThat(e).hasMessageContaining("expected 1, actual 2");
}
StepVerifier.create(mono).expectError(IncorrectResultSizeDataAccessException.class).verify();
}
@Test // DATACASS-335
@@ -741,7 +714,7 @@ public class ReactiveCqlTemplateUnitTests {
Mono<String> mono = template.queryForObject("SELECT * FROM user WHERE username = ?", String.class, "Walter");
assertThat(mono.block()).isEqualTo("OK");
StepVerifier.create(mono).expectNext("OK").verifyComplete();
}
@Test // DATACASS-335
@@ -757,7 +730,7 @@ public class ReactiveCqlTemplateUnitTests {
Flux<String> flux = template.queryForFlux("SELECT * FROM user WHERE username = ?", String.class, "Walter");
assertThat(flux.collectList().block()).contains("OK", "NOT OK");
StepVerifier.create(flux).expectNext("OK", "NOT OK").verifyComplete();
}
@Test // DATACASS-335
@@ -770,7 +743,7 @@ public class ReactiveCqlTemplateUnitTests {
Flux<Row> flux = template.queryForRows("SELECT * FROM user WHERE username = ?", "Walter");
assertThat(flux.collectList().block()).hasSize(2).contains(row);
StepVerifier.create(flux).expectNextCount(2).verifyComplete();
}
@Test // DATACASS-335
@@ -783,7 +756,7 @@ public class ReactiveCqlTemplateUnitTests {
Mono<Boolean> mono = template.execute("UPDATE user SET username = ?", "Walter");
assertThat(mono.block()).isTrue();
StepVerifier.create(mono).expectNext(true).verifyComplete();
}
@Test // DATACASS-335
@@ -798,7 +771,8 @@ public class ReactiveCqlTemplateUnitTests {
Flux<Boolean> flux = template.execute("UPDATE user SET username = ?",
Flux.just(new Object[] { "Walter" }, new Object[] { "Hank" }));
assertThat(flux.collectList().block()).hasSize(2).contains(true);
StepVerifier.create(flux).expectNext(true, true).verifyComplete();
verify(session, atMost(1)).prepare("UPDATE user SET username = ?");
verify(session, times(2)).execute(boundStatement);
}

View File

@@ -79,6 +79,13 @@
<optional>true</optional>
</dependency>
<dependency>
<groupId>io.projectreactor.addons</groupId>
<artifactId>reactor-test</artifactId>
<version>${reactor}</version>
<scope>test</scope>
</dependency>
<dependency>
<groupId>io.reactivex</groupId>
<artifactId>rxjava</artifactId>

View File

@@ -15,10 +15,9 @@
*/
package org.springframework.data.cassandra.core;
import static org.assertj.core.api.Assertions.*;
import reactor.core.publisher.Mono;
import reactor.core.scheduler.Schedulers;
import reactor.test.StepVerifier;
import org.junit.Before;
import org.junit.Test;
@@ -36,7 +35,7 @@ import org.springframework.data.cassandra.test.integration.support.SchemaTestUti
*/
public class ReactiveCassandraTemplateIntegrationTests extends AbstractKeyspaceCreatingIntegrationTest {
private ReactiveCassandraTemplate template;
ReactiveCassandraTemplate template;
@Before
public void setUp() throws Exception {
@@ -57,14 +56,11 @@ public class ReactiveCassandraTemplateIntegrationTests extends AbstractKeyspaceC
Person person = new Person("heisenberg", "Walter", "White");
Mono<Person> insert = template.insert(person);
Mono<Person> oneById = template.selectOneById(person.getId(), Person.class);
StepVerifier.create(template.selectOneById(person.getId(), Person.class)).verifyComplete();
assertThat(oneById.hasElement().block()).isFalse();
StepVerifier.create(insert).expectNext(person).verifyComplete();
Person saved = insert.block();
assertThat(saved).isNotNull().isEqualTo(person);
assertThat(oneById.block()).isNotNull().isEqualTo(saved);
StepVerifier.create(template.selectOneById(person.getId(), Person.class)).expectNext(person).verifyComplete();
}
@Test // DATACASS-335
@@ -72,11 +68,9 @@ public class ReactiveCassandraTemplateIntegrationTests extends AbstractKeyspaceC
Person person = new Person("heisenberg", "Walter", "White");
template.insert(person).block();
StepVerifier.create(template.insert(person)).expectNextCount(1).verifyComplete();
Mono<Long> count = template.count(Person.class);
assertThat(count.block()).isEqualTo(1L);
StepVerifier.create(template.count(Person.class)).expectNext(1L).verifyComplete();
}
@Test // DATACASS-335
@@ -84,17 +78,13 @@ public class ReactiveCassandraTemplateIntegrationTests extends AbstractKeyspaceC
Person person = new Person("heisenberg", "Walter", "White");
template.insert(person).block();
StepVerifier.create(template.insert(person)).expectNextCount(1).verifyComplete();
person.setFirstname("Walter Hartwell");
Person updated = template.update(person).block();
StepVerifier.create(template.insert(person)).expectNextCount(1).verifyComplete();
assertThat(updated).isNotNull();
Mono<Person> oneById = template.selectOneById(person.getId(), Person.class);
assertThat(oneById.block()).isEqualTo(person);
StepVerifier.create(template.selectOneById(person.getId(), Person.class)).expectNext(person).verifyComplete();
}
@Test // DATACASS-335
@@ -102,15 +92,11 @@ public class ReactiveCassandraTemplateIntegrationTests extends AbstractKeyspaceC
Person person = new Person("heisenberg", "Walter", "White");
template.insert(person).block();
StepVerifier.create(template.insert(person)).expectNextCount(1).verifyComplete();
Person deleted = template.delete(person).block();
StepVerifier.create(template.delete(person)).expectNext(person).verifyComplete();
assertThat(deleted).isNotNull();
Mono<Person> oneById = template.selectOneById(person.getId(), Person.class);
assertThat(oneById.block()).isNull();
StepVerifier.create(template.selectOneById(person.getId(), Person.class)).verifyComplete();
}
@Test // DATACASS-335
@@ -118,14 +104,10 @@ public class ReactiveCassandraTemplateIntegrationTests extends AbstractKeyspaceC
Person person = new Person("heisenberg", "Walter", "White");
template.insert(person).block();
StepVerifier.create(template.insert(person)).expectNextCount(1).verifyComplete();
Boolean deleted = template.deleteById(person.getId(), Person.class).block();
StepVerifier.create(template.deleteById(person.getId(), Person.class)).expectNext(true).verifyComplete();
assertThat(deleted).isTrue();
Mono<Person> oneById = template.selectOneById(person.getId(), Person.class);
assertThat(oneById.block()).isNull();
StepVerifier.create(template.selectOneById(person.getId(), Person.class)).verifyComplete();
}
}

View File

@@ -23,6 +23,7 @@ import static org.mockito.Mockito.*;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import reactor.test.StepVerifier;
import java.util.Collections;
@@ -35,7 +36,6 @@ import org.mockito.Mock;
import org.mockito.junit.MockitoJUnitRunner;
import org.springframework.cassandra.core.session.ReactiveResultSet;
import org.springframework.cassandra.core.session.ReactiveSession;
import org.springframework.cassandra.support.exception.CassandraConnectionFailureException;
import org.springframework.data.cassandra.domain.Person;
import com.datastax.driver.core.ColumnDefinitions;
@@ -85,9 +85,10 @@ public class ReactiveCassandraTemplateUnitTests {
when(row.getObject(1)).thenReturn("Walter");
when(row.getObject(2)).thenReturn("White");
Flux<Person> flux = template.select("SELECT * FROM person", Person.class);
StepVerifier.create(template.select("SELECT * FROM person", Person.class)) //
.expectNext(new Person("myid", "Walter", "White")) //
.verifyComplete();
assertThat(flux.collectList().block()).hasSize(1).contains(new Person("myid", "Walter", "White"));
verify(session).execute(statementCaptor.capture());
assertThat(statementCaptor.getValue().toString()).isEqualTo("SELECT * FROM person");
}
@@ -97,15 +98,10 @@ public class ReactiveCassandraTemplateUnitTests {
when(reactiveResultSet.rows()).thenThrow(new NoHostAvailableException(Collections.emptyMap()));
Flux<Person> flux = template.select("SELECT * FROM person", Person.class);
try {
flux.last().block();
fail("Missing CassandraConnectionFailureException");
} catch (CassandraConnectionFailureException e) {
assertThat(e).hasRootCauseInstanceOf(NoHostAvailableException.class);
}
StepVerifier.create(template.select("SELECT * FROM person", Person.class)) //
.consumeErrorWith(e -> {
assertThat(e).hasRootCauseInstanceOf(NoHostAvailableException.class);
}).verify();
}
@Test // DATACASS-335
@@ -123,9 +119,10 @@ public class ReactiveCassandraTemplateUnitTests {
when(row.getObject(1)).thenReturn("Walter");
when(row.getObject(2)).thenReturn("White");
Mono<Person> mono = template.selectOneById("myid", Person.class);
StepVerifier.create(template.selectOneById("myid", Person.class)) //
.expectNext(new Person("myid", "Walter", "White")) //
.verifyComplete();
assertThat(mono.block()).isEqualTo(new Person("myid", "Walter", "White"));
verify(session).execute(statementCaptor.capture());
assertThat(statementCaptor.getValue().toString()).isEqualTo("SELECT * FROM person WHERE id='myid';");
}
@@ -135,9 +132,8 @@ public class ReactiveCassandraTemplateUnitTests {
when(reactiveResultSet.rows()).thenReturn(Flux.just(row));
Mono<Boolean> mono = template.exists("myid", Person.class);
StepVerifier.create(template.exists("myid", Person.class)).expectNext(true).verifyComplete();
assertThat(mono.block()).isTrue();
verify(session).execute(statementCaptor.capture());
assertThat(statementCaptor.getValue().toString()).isEqualTo("SELECT * FROM person WHERE id='myid';");
}
@@ -147,9 +143,8 @@ public class ReactiveCassandraTemplateUnitTests {
when(reactiveResultSet.rows()).thenReturn(Flux.empty());
Mono<Boolean> mono = template.exists("myid", Person.class);
StepVerifier.create(template.exists("myid", Person.class)).expectNext(false).verifyComplete();
assertThat(mono.block()).isFalse();
verify(session).execute(statementCaptor.capture());
assertThat(statementCaptor.getValue().toString()).isEqualTo("SELECT * FROM person WHERE id='myid';");
}
@@ -161,9 +156,8 @@ public class ReactiveCassandraTemplateUnitTests {
when(row.getLong(0)).thenReturn(42L);
when(columnDefinitions.size()).thenReturn(1);
Mono<Long> mono = template.count(Person.class);
StepVerifier.create(template.count(Person.class)).expectNext(42L).verifyComplete();
assertThat(mono.block()).isEqualTo(42L);
verify(session).execute(statementCaptor.capture());
assertThat(statementCaptor.getValue().toString()).isEqualTo("SELECT count(*) FROM person;");
}
@@ -174,9 +168,8 @@ public class ReactiveCassandraTemplateUnitTests {
when(reactiveResultSet.wasApplied()).thenReturn(true);
Person person = new Person("heisenberg", "Walter", "White");
Mono<Person> mono = template.insert(person);
StepVerifier.create(template.insert(person)).expectNext(person).verifyComplete();
assertThat(mono.block()).isEqualTo(person);
verify(session).execute(statementCaptor.capture());
assertThat(statementCaptor.getValue().toString())
.isEqualTo("INSERT INTO person (firstname,id,lastname) VALUES ('Walter','heisenberg','White');");
@@ -189,15 +182,11 @@ public class ReactiveCassandraTemplateUnitTests {
when(session.execute(any(Statement.class)))
.thenReturn(Mono.error(new NoHostAvailableException(Collections.emptyMap())));
Mono<Person> mono = template.insert(new Person("heisenberg", "Walter", "White"));
StepVerifier.create(template.insert(new Person("heisenberg", "Walter", "White"))) //
.consumeErrorWith(e -> {
try {
mono.block();
fail("Missing CassandraConnectionFailureException");
} catch (CassandraConnectionFailureException e) {
assertThat(e).hasRootCauseInstanceOf(NoHostAvailableException.class);
}
assertThat(e).hasRootCauseInstanceOf(NoHostAvailableException.class);
}).verify();
}
@Test // DATACASS-335
@@ -206,9 +195,8 @@ public class ReactiveCassandraTemplateUnitTests {
when(reactiveResultSet.wasApplied()).thenReturn(false);
Person person = new Person("heisenberg", "Walter", "White");
Mono<Person> mono = template.insert(person);
assertThat(mono.block()).isNull();
StepVerifier.create(template.insert(person)).verifyComplete();
}
@Test // DATACASS-335
@@ -217,41 +205,22 @@ public class ReactiveCassandraTemplateUnitTests {
when(reactiveResultSet.wasApplied()).thenReturn(true);
Person person = new Person("heisenberg", "Walter", "White");
Mono<Person> mono = template.update(person);
assertThat(mono.block()).isEqualTo(person);
StepVerifier.create(template.update(person)).expectNext(person).verifyComplete();
verify(session).execute(statementCaptor.capture());
assertThat(statementCaptor.getValue().toString())
.isEqualTo("UPDATE person SET firstname='Walter',lastname='White' WHERE id='heisenberg';");
}
@Test // DATACASS-335
public void updateShouldTranslateException() {
reset(session);
when(session.execute(any(Statement.class)))
.thenReturn(Mono.error(new NoHostAvailableException(Collections.emptyMap())));
Mono<Person> mono = template.update(new Person("heisenberg", "Walter", "White"));
try {
mono.block();
fail("Missing CassandraConnectionFailureException");
} catch (CassandraConnectionFailureException e) {
assertThat(e).hasRootCauseInstanceOf(NoHostAvailableException.class);
}
}
@Test // DATACASS-335
public void updateShouldNotApplyUpdate() {
when(reactiveResultSet.wasApplied()).thenReturn(false);
Person person = new Person("heisenberg", "Walter", "White");
Mono<Person> mono = template.update(person);
assertThat(mono.block()).isNull();
StepVerifier.create(template.update(person)).verifyComplete();
}
@Test // DATACASS-335
@@ -261,46 +230,26 @@ public class ReactiveCassandraTemplateUnitTests {
Person person = new Person("heisenberg", "Walter", "White");
Mono<Person> mono = template.delete(person);
StepVerifier.create(template.delete(person)).expectNext(person).verifyComplete();
assertThat(mono.block()).isEqualTo(person);
verify(session).execute(statementCaptor.capture());
assertThat(statementCaptor.getValue().toString()).isEqualTo("DELETE FROM person WHERE id='heisenberg';");
}
@Test // DATACASS-335
public void deleteShouldTranslateException() {
reset(session);
when(session.execute(any(Statement.class)))
.thenReturn(Mono.error(new NoHostAvailableException(Collections.emptyMap())));
Mono<Person> mono = template.delete(new Person("heisenberg", "Walter", "White"));
try {
mono.block();
fail("Missing CassandraConnectionFailureException");
} catch (CassandraConnectionFailureException e) {
assertThat(e).hasRootCauseInstanceOf(NoHostAvailableException.class);
}
}
@Test // DATACASS-335
public void deleteShouldNotApplyRemoval() {
when(reactiveResultSet.wasApplied()).thenReturn(false);
Person person = new Person("heisenberg", "Walter", "White");
Mono<Person> mono = template.delete(person);
assertThat(mono.block()).isNull();
StepVerifier.create(template.delete(person)).verifyComplete();
}
@Test // DATACASS-335
public void truncateShouldRemoveEntities() {
template.truncate(Person.class).block();
StepVerifier.create(template.truncate(Person.class)).verifyComplete();
verify(session).execute(statementCaptor.capture());
assertThat(statementCaptor.getValue().toString()).isEqualTo("TRUNCATE person;");

View File

@@ -19,11 +19,12 @@ import static org.assertj.core.api.Assertions.*;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import reactor.test.TestSubscriber;
import reactor.test.StepVerifier;
import rx.Observable;
import rx.Single;
import java.util.Arrays;
import java.util.List;
import org.junit.Before;
import org.junit.Test;
@@ -33,7 +34,6 @@ import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.cassandra.test.integration.AbstractKeyspaceCreatingIntegrationTest;
import org.springframework.context.annotation.ComponentScan.Filter;
import org.springframework.context.annotation.Configuration;
import org.springframework.data.cassandra.core.ReactiveCassandraTemplate;
import org.springframework.data.cassandra.domain.Person;
import org.springframework.data.cassandra.repository.config.EnableReactiveCassandraRepositories;
import org.springframework.data.cassandra.test.integration.support.IntegrationTestConfig;
@@ -68,7 +68,6 @@ public class ConvertingReactiveCassandraRepositoryTests extends AbstractKeyspace
}
@Autowired Session session;
@Autowired ReactiveCassandraTemplate template;
@Autowired MixedPersonRepostitory reactiveRepository;
@Autowired PersonRepostitory reactivePersonRepostitory;
@Autowired RxJava1PersonRepostitory rxJava1PersonRepostitory;
@@ -86,130 +85,114 @@ public class ConvertingReactiveCassandraRepositoryTests extends AbstractKeyspace
Thread.sleep(500);
}
reactiveRepository.deleteAll().block();
StepVerifier.create(reactiveRepository.deleteAll()).verifyComplete();
dave = new Person("42", "Dave", "Matthews");
oliver = new Person("4", "Oliver August", "Matthews");
carter = new Person("49", "Carter", "Beauford");
boyd = new Person("45", "Boyd", "Tinsley");
TestSubscriber<Person> subscriber = TestSubscriber.create();
reactiveRepository.save(Arrays.asList(oliver, dave, carter, boyd)).subscribe(subscriber);
subscriber.await().assertComplete().assertNoError();
StepVerifier.create(reactiveRepository.save(Arrays.asList(oliver, dave, carter, boyd))).expectNextCount(4)
.verifyComplete();
}
@Test // DATACASS-335
public void reactiveStreamsMethodsShouldWork() throws InterruptedException {
TestSubscriber<Boolean> subscriber = TestSubscriber.subscribe(reactivePersonRepostitory.exists(dave.getId()));
subscriber.awaitAndAssertNextValueCount(1).assertNoError().assertValues(true);
public void reactiveStreamsMethodsShouldWork() {
StepVerifier.create(reactivePersonRepostitory.exists(dave.getId())).expectNext(true).verifyComplete();
}
@Test // DATACASS-335
public void reactiveStreamsQueryMethodsShouldWork() {
TestSubscriber<Person> subscriber = TestSubscriber
.subscribe(reactivePersonRepostitory.findByLastname(boyd.getLastname()));
subscriber.awaitAndAssertNextValueCount(1).assertValues(boyd);
StepVerifier.create(reactivePersonRepostitory.findByLastname(boyd.getLastname())).expectNext(boyd).verifyComplete();
}
@Test // DATACASS-360
public void dtoProjectionShouldWork() {
TestSubscriber<PersonDto> subscriber = TestSubscriber
.subscribe(reactivePersonRepostitory.findProjectedByLastname(boyd.getLastname()));
StepVerifier.create(reactivePersonRepostitory.findProjectedByLastname(boyd.getLastname()))
.consumeNextWith(actual -> {
subscriber.awaitAndAssertNextValueCount(1).assertValuesWith(personDto -> {
assertThat(personDto.firstname).isEqualTo(boyd.getFirstname());
assertThat(personDto.lastname).isEqualTo(boyd.getLastname());
});
assertThat(actual.firstname).isEqualTo(boyd.getFirstname());
assertThat(actual.lastname).isEqualTo(boyd.getLastname());
}).verifyComplete();
}
@Test // DATACASS-335
public void simpleRxJavaMethodsShouldWork() {
rx.observers.TestSubscriber<Boolean> subscriber = new rx.observers.TestSubscriber<>();
rxJava1PersonRepostitory.exists(dave.getId()).subscribe(subscriber);
subscriber.awaitTerminalEvent();
subscriber.assertCompleted();
subscriber.assertNoErrors();
subscriber.assertValue(true);
rxJava1PersonRepostitory.exists(dave.getId()) //
.test() //
.awaitTerminalEvent() //
.assertResult(true) //
.assertCompleted() //
.assertNoErrors();
}
@Test // DATACASS-335
public void existsWithSingleRxJavaIdMethodsShouldWork() {
rx.observers.TestSubscriber<Boolean> subscriber = new rx.observers.TestSubscriber<>();
rxJava1PersonRepostitory.exists(Single.just(dave.getId())).subscribe(subscriber);
subscriber.awaitTerminalEvent();
subscriber.assertCompleted();
subscriber.assertNoErrors();
subscriber.assertValue(true);
rxJava1PersonRepostitory.exists(Single.just(dave.getId())) //
.test() //
.awaitTerminalEvent() //
.assertResult(true) //
.assertCompleted() //
.assertNoErrors();
}
@Test // DATACASS-335
public void singleRxJavaQueryMethodShouldWork() {
rx.observers.TestSubscriber<Person> subscriber = new rx.observers.TestSubscriber<>();
rxJava1PersonRepostitory.findManyByLastname(dave.getLastname()).subscribe(subscriber);
subscriber.awaitTerminalEvent();
subscriber.assertNoErrors();
subscriber.assertCompleted();
subscriber.assertValueCount(2);
rxJava1PersonRepostitory.findManyByLastname(dave.getLastname()) //
.test() //
.awaitTerminalEvent() //
.assertValueCount(2) //
.assertNoErrors() //
.assertCompleted();
}
@Test // DATACASS-335
public void singleProjectedRxJavaQueryMethodShouldWork() {
rx.observers.TestSubscriber<ProjectedPerson> subscriber = new rx.observers.TestSubscriber<>();
List<ProjectedPerson> values = rxJava1PersonRepostitory.findProjectedByLastname(carter.getLastname()) //
.test() //
.awaitTerminalEvent() //
.assertValueCount(1) //
.assertCompleted() //
.assertNoErrors() //
.getOnNextEvents();
rxJava1PersonRepostitory.findProjectedByLastname(carter.getLastname()).subscribe(subscriber);
subscriber.awaitTerminalEvent();
subscriber.assertCompleted();
subscriber.assertNoErrors();
ProjectedPerson projectedPerson = subscriber.getOnNextEvents().get(0);
ProjectedPerson projectedPerson = values.get(0);
assertThat(projectedPerson.getFirstname()).isEqualTo(carter.getFirstname());
}
@Test // DATACASS-335
public void observableRxJavaQueryMethodShouldWork() {
rx.observers.TestSubscriber<Person> subscriber = new rx.observers.TestSubscriber<>();
rxJava1PersonRepostitory.findByLastname(boyd.getLastname()).subscribe(subscriber);
subscriber.awaitTerminalEvent();
subscriber.assertCompleted();
subscriber.assertNoErrors();
subscriber.assertValue(boyd);
rxJava1PersonRepostitory.findByLastname(boyd.getLastname()) //
.test() //
.awaitTerminalEvent() //
.assertValue(boyd) //
.assertNoErrors() //
.assertCompleted();
}
@Test // DATACASS-335
public void mixedRepositoryShouldWork() {
Person value = reactiveRepository.findByLastname(boyd.getLastname()).toBlocking().value();
assertThat(value).isEqualTo(boyd);
reactiveRepository.findByLastname(boyd.getLastname()) //
.test() //
.awaitTerminalEvent() //
.assertValue(boyd) //
.assertCompleted() //
.assertNoErrors();
}
@Test // DATACASS-335
public void shouldFindOneByPublisherOfLastName() {
Person carter = reactiveRepository.findByLastname(Single.just(this.carter.getLastname())).block();
StepVerifier.create(reactiveRepository.findByLastname(Single.just(this.carter.getLastname()))) //
.expectNext(carter) //
.verifyComplete();
assertThat(carter.getFirstname()).isEqualTo(this.carter.getFirstname());
}
@Repository

View File

@@ -15,13 +15,11 @@
*/
package org.springframework.data.cassandra.repository;
import static org.assertj.core.api.Assertions.*;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import reactor.test.StepVerifier;
import java.util.Arrays;
import java.util.List;
import org.junit.Before;
import org.junit.Test;
@@ -111,47 +109,34 @@ public class ReactiveCassandraRepositoryIntegrationTests extends AbstractKeyspac
repository = factory.getRepository(PersonRepository.class);
groupRepostitory = factory.getRepository(GroupRepository.class);
repository.deleteAll().block();
groupRepostitory.deleteAll().block();
StepVerifier.create(repository.deleteAll().concatWith(groupRepostitory.deleteAll())).verifyComplete();
dave = new Person("42", "Dave", "Matthews");
oliver = new Person("4", "Oliver August", "Matthews");
carter = new Person("49", "Carter", "Beauford");
boyd = new Person("45", "Boyd", "Tinsley");
repository.save(Arrays.asList(oliver, dave, carter, boyd)).last().block();
StepVerifier.create(repository.save(Arrays.asList(oliver, dave, carter, boyd))).expectNextCount(4).verifyComplete();
}
@Test // DATACASS-335
public void shouldFindByLastName() {
List<Person> list = repository.findByLastname("Matthews").collectList().block();
assertThat(list).hasSize(2).contains(dave, oliver);
StepVerifier.create(repository.findByLastname(dave.getLastname())).expectNextCount(2).verifyComplete();
}
@Test // DATACASS-335
public void shouldFindOneByLastName() {
Person carter = repository.findOneByLastname("Beauford").block();
assertThat(carter.getFirstname()).isEqualTo("Carter");
StepVerifier.create(repository.findOneByLastname(carter.getLastname())).expectNext(carter).verifyComplete();
}
@Test // DATACASS-335
public void shouldFindOneByPublisherOfLastName() {
Person carter = repository.findByLastname(Mono.just("Beauford")).block();
assertThat(carter.getFirstname()).isEqualTo("Carter");
StepVerifier.create(repository.findByLastname(Mono.just(carter.getLastname()))).expectNext(carter).verifyComplete();
}
@Test // DATACASS-335
public void shouldFindUsingPublishersInStringQuery() {
List<Person> persons = repository.findStringQuery(Mono.just("Matthews")).collectList().block();
assertThat(persons).contains(dave);
StepVerifier.create(repository.findStringQuery(Mono.just(dave.getLastname()))).expectNextCount(2).verifyComplete();
}
@Test // DATACASS-335
@@ -160,17 +145,20 @@ public class ReactiveCassandraRepositoryIntegrationTests extends AbstractKeyspac
GroupKey key1 = new GroupKey("Simpsons", "hash", "Bart");
GroupKey key2 = new GroupKey("Simpsons", "hash", "Homer");
groupRepostitory.save(Flux.just(new Group(key1), new Group(key2))).blockLast();
StepVerifier.create(groupRepostitory.save(Flux.just(new Group(key1), new Group(key2)))).expectNextCount(2)
.verifyComplete();
List<Group> persons = groupRepostitory
.findByIdGroupnameAndIdHashPrefix("Simpsons", "hash", new Sort(Direction.ASC, "id.username")).collectList()
.block();
assertThat(persons).containsSequence(new Group(key1), new Group(key2));
StepVerifier
.create(groupRepostitory.findByIdGroupnameAndIdHashPrefix("Simpsons", "hash",
new Sort(Direction.ASC, "id.username"))) //
.expectNext(new Group(key1), new Group(key2)) //
.verifyComplete();
List<Group> reversed = groupRepostitory
.findByIdGroupnameAndIdHashPrefix("Simpsons", "hash", new Sort(Direction.DESC, "id.username")).collectList()
.block();
assertThat(reversed).containsSequence(new Group(key2), new Group(key1));
StepVerifier
.create(groupRepostitory.findByIdGroupnameAndIdHashPrefix("Simpsons", "hash",
new Sort(Direction.DESC, "id.username"))) //
.expectNext(new Group(key2), new Group(key1)) //
.verifyComplete();
}
interface PersonRepository extends ReactiveCassandraRepository<Person, String> {

View File

@@ -19,10 +19,9 @@ import static org.assertj.core.api.Assertions.*;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import reactor.test.TestSubscriber;
import reactor.test.StepVerifier;
import java.util.Arrays;
import java.util.List;
import org.junit.Before;
import org.junit.Test;
@@ -91,166 +90,149 @@ public class SimpleReactiveCassandraRepositoryIntegrationTests extends AbstractK
repository = factory.getRepository(PersonRepostitory.class);
repository.deleteAll().block();
deleteAll();
dave = new Person("42", "Dave", "Matthews");
oliver = new Person("4", "Oliver August", "Matthews");
carter = new Person("49", "Carter", "Beauford");
boyd = new Person("45", "Boyd", "Tinsley");
}
repository.save(Arrays.asList(oliver, dave, carter, boyd)).last().block();
private void insertTestData() {
StepVerifier.create(repository.save(Arrays.asList(oliver, dave, carter, boyd))).expectNextCount(4).verifyComplete();
}
private void deleteAll() {
StepVerifier.create(repository.deleteAll()).verifyComplete();
}
@Test // DATACASS-335
public void existsByIdShouldReturnTrueForExistingObject() {
Boolean exists = repository.exists(dave.getId()).block();
insertTestData();
assertThat(exists).isTrue();
StepVerifier.create(repository.exists(dave.getId())).expectNext(true).verifyComplete();
}
@Test // DATACASS-335
public void existsByIdShouldReturnFalseForAbsentObject() {
TestSubscriber<Boolean> testSubscriber = TestSubscriber.subscribe(repository.exists("unknown"));
testSubscriber.await().assertComplete().assertValues(false).assertNoError();
StepVerifier.create(repository.exists("unknown")).expectNext(false).verifyComplete();
}
@Test // DATACASS-335
public void existsByMonoOfIdShouldReturnTrueForExistingObject() {
Boolean exists = repository.exists(Mono.just(dave.getId())).block();
assertThat(exists).isTrue();
insertTestData();
StepVerifier.create(repository.exists(Mono.just(dave.getId()))).expectNext(true).verifyComplete();
}
@Test // DATACASS-335
public void existsByEmptyMonoOfIdShouldReturnEmptyMono() {
TestSubscriber<Boolean> testSubscriber = TestSubscriber.subscribe(repository.exists(Mono.empty()));
testSubscriber.await().assertComplete().assertNoValues().assertNoError();
StepVerifier.create(repository.exists(Mono.empty())).verifyComplete();
}
@Test // DATACASS-335
public void findOneShouldReturnObject() {
Person person = repository.findOne(dave.getId()).block();
insertTestData();
assertThat(person).isEqualTo(dave);
StepVerifier.create(repository.findOne(dave.getId())).expectNext(dave).verifyComplete();
}
@Test // DATACASS-335
public void findOneShouldCompleteWithoutValueForAbsentObject() {
TestSubscriber<Person> testSubscriber = TestSubscriber.subscribe(repository.findOne("unknown"));
testSubscriber.await().assertComplete().assertNoValues().assertNoError();
StepVerifier.create(repository.findOne("unknown")).verifyComplete();
}
@Test // DATACASS-335
public void findOneByMonoOfIdShouldReturnTrueForExistingObject() {
Person person = repository.findOne(Mono.just(dave.getId())).block();
insertTestData();
assertThat(person).isEqualTo(dave);
StepVerifier.create(repository.findOne(Mono.just(dave.getId()))).expectNext(dave).verifyComplete();
}
@Test // DATACASS-335
public void findOneByEmptyMonoOfIdShouldReturnEmptyMono() {
TestSubscriber<Person> testSubscriber = TestSubscriber.subscribe(repository.findOne(Mono.empty()));
testSubscriber.await().assertComplete().assertNoValues().assertNoError();
StepVerifier.create(repository.findOne(Mono.empty())).verifyComplete();
}
@Test // DATACASS-335
public void findAllShouldReturnAllResults() {
List<Person> persons = repository.findAll().collectList().block();
insertTestData();
assertThat(persons).hasSize(4);
StepVerifier.create(repository.findAll()).expectNextCount(4).verifyComplete();
}
@Test // DATACASS-335
public void findAllByIterableOfIdShouldReturnResults() {
List<Person> persons = repository.findAll(Arrays.asList(dave.getId(), boyd.getId())).collectList().block();
insertTestData();
assertThat(persons).hasSize(2);
StepVerifier.create(repository.findAll(Arrays.asList(dave.getId(), boyd.getId()))) //
.expectNextCount(2) //
.verifyComplete();
}
@Test // DATACASS-335
public void findAllByPublisherOfIdShouldReturnResults() {
List<Person> persons = repository.findAll(Flux.just(dave.getId(), boyd.getId())).collectList().block();
insertTestData();
assertThat(persons).hasSize(2);
StepVerifier.create(repository.findAll(Flux.just(dave.getId(), boyd.getId()))) //
.expectNextCount(2) //
.verifyComplete();
}
@Test // DATACASS-335
public void findAllByEmptyPublisherOfIdShouldReturnResults() {
TestSubscriber<Person> testSubscriber = TestSubscriber.subscribe(repository.findAll(Flux.empty()));
testSubscriber.await().assertComplete().assertNoValues().assertNoError();
StepVerifier.create(repository.findAll(Flux.empty())).verifyComplete();
}
@Test // DATACASS-335
public void countShouldReturnNumberOfRecords() {
TestSubscriber<Long> testSubscriber = TestSubscriber.subscribe(repository.count());
insertTestData();
testSubscriber.await().assertComplete().assertValueCount(1).assertValues(4L).assertNoError();
StepVerifier.create(repository.count()).expectNext(4L).verifyComplete();
}
@Test // DATACASS-335
public void insertEntityShouldInsertEntity() {
repository.deleteAll().block();
Person person = new Person("36", "Homer", "Simpson");
TestSubscriber<Person> testSubscriber = TestSubscriber.subscribe(repository.insert(person));
StepVerifier.create(repository.insert(person)).expectNext(person).verifyComplete();
testSubscriber.await().assertComplete().assertValueCount(1).assertValues(person);
repository.findAll().count().subscribeWith(TestSubscriber.create()).awaitAndAssertNextValues(1L);
StepVerifier.create(repository.findAll()).expectNextCount(1L).verifyComplete();
}
@Test // DATACASS-335
public void insertShouldDeferredWrite() {
repository.deleteAll().block();
Person person = new Person("36", "Homer", "Simpson");
repository.insert(person);
repository.findAll().count().subscribeWith(TestSubscriber.create()).awaitAndAssertNextValues(0L);
StepVerifier.create(repository.findAll()).expectNextCount(0L).verifyComplete();
}
@Test // DATACASS-335
public void insertIterableOfEntitiesShouldInsertEntity() {
repository.deleteAll().block();
StepVerifier.create(repository.insert(Arrays.asList(dave, oliver, boyd))).expectNextCount(3L).verifyComplete();
TestSubscriber<Person> testSubscriber = TestSubscriber
.subscribe(repository.insert(Arrays.asList(dave, oliver, boyd)));
testSubscriber.await().assertComplete().assertValueCount(3);
repository.findAll().count().subscribeWith(TestSubscriber.create()).awaitAndAssertNextValues(3L);
StepVerifier.create(repository.findAll()).expectNextCount(3L).verifyComplete();
}
@Test // DATACASS-335
public void insertPublisherOfEntitiesShouldInsertEntity() {
repository.deleteAll().block();
StepVerifier.create(repository.insert(Flux.just(dave, oliver, boyd))).expectNextCount(3L).verifyComplete();
TestSubscriber<Person> testSubscriber = TestSubscriber.subscribe(repository.insert(Flux.just(dave, oliver, boyd)));
testSubscriber.await().assertComplete().assertValueCount(3);
repository.findAll().count().subscribeWith(TestSubscriber.create()).awaitAndAssertNextValues(3L);
StepVerifier.create(repository.findAll()).expectNextCount(3L).verifyComplete();
}
@Test // DATACASS-335
@@ -259,14 +241,13 @@ public class SimpleReactiveCassandraRepositoryIntegrationTests extends AbstractK
dave.setFirstname("Hello, Dave");
dave.setLastname("Bowman");
TestSubscriber<Person> testSubscriber = TestSubscriber.subscribe(repository.save(dave));
StepVerifier.create(repository.save(dave)).expectNextCount(1).verifyComplete();
testSubscriber.await().assertComplete().assertValueCount(1).assertValues(dave);
StepVerifier.create(repository.findOne(dave.getId())).consumeNextWith(actual -> {
Person loaded = repository.findOne(dave.getId()).block();
assertThat(loaded.getFirstname()).isEqualTo(dave.getFirstname());
assertThat(loaded.getLastname()).isEqualTo(dave.getLastname());
assertThat(actual.getFirstname()).isEqualTo(dave.getFirstname());
assertThat(actual.getLastname()).isEqualTo(dave.getLastname());
}).verifyComplete();
}
@Test // DATACASS-335
@@ -274,26 +255,17 @@ public class SimpleReactiveCassandraRepositoryIntegrationTests extends AbstractK
Person person = new Person("36", "Homer", "Simpson");
TestSubscriber<Person> testSubscriber = TestSubscriber.subscribe(repository.save(person));
StepVerifier.create(repository.save(person)).expectNextCount(1).verifyComplete();
testSubscriber.await().assertComplete().assertValueCount(1).assertValues(person);
Person loaded = repository.findOne(person.getId()).block();
assertThat(loaded).isEqualTo(person);
StepVerifier.create(repository.findOne(person.getId())).expectNext(person).verifyComplete();
}
@Test // DATACASS-335
public void saveIterableOfNewEntitiesShouldInsertEntity() {
repository.deleteAll().block();
StepVerifier.create(repository.save(Arrays.asList(dave, oliver, boyd))).expectNextCount(3).verifyComplete();
TestSubscriber<Person> testSubscriber = TestSubscriber
.subscribe(repository.save(Arrays.asList(dave, oliver, boyd)));
testSubscriber.await().assertComplete().assertValueCount(3);
repository.findAll().count().subscribeWith(TestSubscriber.create()).awaitAndAssertNextValues(3L);
StepVerifier.create(repository.findAll()).expectNextCount(3L).verifyComplete();
}
@Test // DATACASS-335
@@ -304,82 +276,61 @@ public class SimpleReactiveCassandraRepositoryIntegrationTests extends AbstractK
dave.setFirstname("Hello, Dave");
dave.setLastname("Bowman");
TestSubscriber<Person> testSubscriber = TestSubscriber.subscribe(repository.save(Arrays.asList(person, dave)));
StepVerifier.create(repository.save(Arrays.asList(person, dave))).expectNextCount(2).verifyComplete();
testSubscriber.await().assertComplete().assertValueCount(2);
StepVerifier.create(repository.findOne(dave.getId())).expectNext(dave).verifyComplete();
Person persistentDave = repository.findOne(dave.getId()).block();
assertThat(persistentDave).isEqualTo(dave);
Person persistentHomer = repository.findOne(person.getId()).block();
assertThat(persistentHomer).isEqualTo(person);
StepVerifier.create(repository.findOne(person.getId())).expectNext(person).verifyComplete();
}
@Test // DATACASS-335
public void savePublisherOfEntitiesShouldInsertEntity() {
repository.deleteAll().block();
StepVerifier.create(repository.save(Flux.just(dave, oliver, boyd))).expectNextCount(3).verifyComplete();
TestSubscriber<Person> testSubscriber = TestSubscriber.subscribe(repository.save(Flux.just(dave, oliver, boyd)));
testSubscriber.await().assertComplete().assertValueCount(3);
repository.findAll().count().subscribeWith(TestSubscriber.create()).awaitAndAssertNextValues(3L);
StepVerifier.create(repository.findAll()).expectNextCount(3L).verifyComplete();
}
@Test // DATACASS-335
public void deleteAllShouldRemoveEntities() {
repository.deleteAll().block();
insertTestData();
TestSubscriber<Person> testSubscriber = TestSubscriber.subscribe(repository.findAll());
StepVerifier.create(repository.deleteAll()).verifyComplete();
testSubscriber.await().assertComplete().assertValueCount(0);
StepVerifier.create(repository.findAll()).verifyComplete();
}
@Test // DATACASS-335
public void deleteByIdShouldRemoveEntity() {
TestSubscriber<Void> testSubscriber = TestSubscriber.subscribe(repository.delete(dave.getId()));
StepVerifier.create(repository.delete(dave.getId())).verifyComplete();
testSubscriber.await().assertComplete().assertNoValues();
TestSubscriber<Person> verificationSubscriber = TestSubscriber.subscribe(repository.findOne(dave.getId()));
verificationSubscriber.await().assertComplete().assertNoValues();
StepVerifier.create(repository.findOne(dave.getId())).expectNextCount(0).verifyComplete();
}
@Test // DATACASS-335
public void deleteShouldRemoveEntity() {
TestSubscriber<Void> testSubscriber = TestSubscriber.subscribe(repository.delete(dave));
StepVerifier.create(repository.delete(dave)).verifyComplete();
testSubscriber.await().assertComplete().assertNoValues();
TestSubscriber<Person> verificationSubscriber = TestSubscriber.subscribe(repository.findOne(dave.getId()));
verificationSubscriber.await().assertComplete().assertNoValues();
StepVerifier.create(repository.findOne(dave.getId())).expectNextCount(0).verifyComplete();
}
@Test // DATACASS-335
public void deleteIterableOfEntitiesShouldRemoveEntities() {
TestSubscriber<Void> testSubscriber = TestSubscriber.subscribe(repository.delete(Arrays.asList(dave, boyd)));
StepVerifier.create(repository.delete(Arrays.asList(dave, boyd))).verifyComplete();
testSubscriber.await().assertComplete().assertNoValues();
TestSubscriber<Person> verificationSubscriber = TestSubscriber.subscribe(repository.findOne(boyd.getId()));
verificationSubscriber.await().assertComplete().assertNoValues();
StepVerifier.create(repository.findOne(boyd.getId())).expectNextCount(0).verifyComplete();
}
@Test // DATACASS-335
public void deletePublisherOfEntitiesShouldRemoveEntities() {
TestSubscriber<Void> testSubscriber = TestSubscriber.subscribe(repository.delete(Flux.just(dave, boyd)));
StepVerifier.create(repository.delete(Flux.just(dave, boyd))).verifyComplete();
testSubscriber.await().assertComplete().assertNoValues();
TestSubscriber<Person> verificationSubscriber = TestSubscriber.subscribe(repository.findOne(boyd.getId()));
verificationSubscriber.await().assertComplete().assertNoValues();
StepVerifier.create(repository.findOne(boyd.getId())).expectNextCount(0).verifyComplete();
}
interface PersonRepostitory extends ReactiveCassandraRepository<Person, String> {}