diff --git a/spring-data-cassandra/pom.xml b/spring-data-cassandra/pom.xml
index 2493a0774..0208a20c7 100644
--- a/spring-data-cassandra/pom.xml
+++ b/spring-data-cassandra/pom.xml
@@ -182,22 +182,34 @@
kotlin-stdlib
true
+
org.jetbrains.kotlin
kotlin-reflect
true
+
- org.jetbrains.kotlin
- kotlin-test-junit
- test
+ org.jetbrains.kotlinx
+ kotlinx-coroutines-core
+ ${kotlin-coroutines}
+ true
+
+
+ org.jetbrains.kotlinx
+ kotlinx-coroutines-reactor
+ ${kotlin-coroutines}
+ true
+
+
io.mockk
mockk
${mockk}
test
+
diff --git a/spring-data-cassandra/src/main/kotlin/org/springframework/data/cassandra/core/ReactiveDeleteOperationExtensions.kt b/spring-data-cassandra/src/main/kotlin/org/springframework/data/cassandra/core/ReactiveDeleteOperationExtensions.kt
index 12bdfa8b0..5e5dbb993 100644
--- a/spring-data-cassandra/src/main/kotlin/org/springframework/data/cassandra/core/ReactiveDeleteOperationExtensions.kt
+++ b/spring-data-cassandra/src/main/kotlin/org/springframework/data/cassandra/core/ReactiveDeleteOperationExtensions.kt
@@ -15,6 +15,7 @@
*/
package org.springframework.data.cassandra.core
+import kotlinx.coroutines.reactive.awaitSingle
import kotlin.reflect.KClass
/**
@@ -35,3 +36,12 @@ fun ReactiveDeleteOperation.delete(entityClass: KClass): ReactiveDe
*/
inline fun ReactiveDeleteOperation.delete(): ReactiveDeleteOperation.ReactiveDelete =
delete(T::class.java)
+
+/**
+ * Coroutines variant of [ReactiveDeleteOperation.TerminatingDelete.all].
+ *
+ * @author Mark Paluch
+ * @since 2.2
+ */
+suspend fun ReactiveDeleteOperation.TerminatingDelete.allAndAwait(): WriteResult =
+ all().awaitSingle()
diff --git a/spring-data-cassandra/src/main/kotlin/org/springframework/data/cassandra/core/ReactiveInsertOperationExtensions.kt b/spring-data-cassandra/src/main/kotlin/org/springframework/data/cassandra/core/ReactiveInsertOperationExtensions.kt
index e8ff95fce..a06565ac1 100644
--- a/spring-data-cassandra/src/main/kotlin/org/springframework/data/cassandra/core/ReactiveInsertOperationExtensions.kt
+++ b/spring-data-cassandra/src/main/kotlin/org/springframework/data/cassandra/core/ReactiveInsertOperationExtensions.kt
@@ -15,6 +15,7 @@
*/
package org.springframework.data.cassandra.core
+import kotlinx.coroutines.reactive.awaitSingle
import kotlin.reflect.KClass
/**
@@ -35,3 +36,12 @@ fun ReactiveInsertOperation.insert(entityClass: KClass): ReactiveIn
*/
inline fun ReactiveInsertOperation.insert(): ReactiveInsertOperation.ReactiveInsert =
insert(T::class.java)
+
+/**
+ * Coroutines variant of [ReactiveInsertOperation.TerminatingInsert.one].
+ *
+ * @author Mark Paluch
+ * @since 2.2
+ */
+suspend inline fun ReactiveInsertOperation.TerminatingInsert.oneAndAwait(o: T): EntityWriteResult =
+ one(o).awaitSingle()
diff --git a/spring-data-cassandra/src/main/kotlin/org/springframework/data/cassandra/core/ReactiveSelectOperationExtensions.kt b/spring-data-cassandra/src/main/kotlin/org/springframework/data/cassandra/core/ReactiveSelectOperationExtensions.kt
index 9ef970ca4..4b74fbde0 100644
--- a/spring-data-cassandra/src/main/kotlin/org/springframework/data/cassandra/core/ReactiveSelectOperationExtensions.kt
+++ b/spring-data-cassandra/src/main/kotlin/org/springframework/data/cassandra/core/ReactiveSelectOperationExtensions.kt
@@ -15,6 +15,8 @@
*/
package org.springframework.data.cassandra.core
+import kotlinx.coroutines.reactive.awaitFirstOrNull
+import kotlinx.coroutines.reactive.awaitSingle
import kotlin.reflect.KClass
/**
@@ -47,3 +49,39 @@ fun ReactiveSelectOperation.SelectWithProjection<*>.asType(resultType:
*/
inline fun ReactiveSelectOperation.SelectWithProjection<*>.asType(): ReactiveSelectOperation.SelectWithQuery =
`as`(T::class.java)
+
+/**
+ * Coroutines variant of [ReactiveSelectOperation.TerminatingSelect.one].
+ *
+ * @author Mark Paluch
+ * @since 2.2
+ */
+suspend inline fun ReactiveSelectOperation.TerminatingSelect.awaitOne(): T? =
+ one().awaitFirstOrNull()
+
+/**
+ * Coroutines variant of [ReactiveSelectOperation.TerminatingSelect.first].
+ *
+ * @author Mark Paluch
+ * @since 2.2
+ */
+suspend inline fun ReactiveSelectOperation.TerminatingSelect.awaitFirst(): T? =
+ first().awaitFirstOrNull()
+
+/**
+ * Coroutines variant of [ReactiveSelectOperation.TerminatingSelect.count].
+ *
+ * @author Mark Paluch
+ * @since 2.2
+ */
+suspend fun ReactiveSelectOperation.TerminatingSelect.awaitCount(): Long =
+ count().awaitSingle()
+
+/**
+ * Coroutines variant of [ReactiveSelectOperation.TerminatingSelect.exists].
+ *
+ * @author Mark Paluch
+ * @since 2.2
+ */
+suspend fun ReactiveSelectOperation.TerminatingSelect.awaitExists(): Boolean =
+ exists().awaitSingle()
diff --git a/spring-data-cassandra/src/main/kotlin/org/springframework/data/cassandra/core/ReactiveUpdateOperationExtensions.kt b/spring-data-cassandra/src/main/kotlin/org/springframework/data/cassandra/core/ReactiveUpdateOperationExtensions.kt
index c9ade639e..8482a6f95 100644
--- a/spring-data-cassandra/src/main/kotlin/org/springframework/data/cassandra/core/ReactiveUpdateOperationExtensions.kt
+++ b/spring-data-cassandra/src/main/kotlin/org/springframework/data/cassandra/core/ReactiveUpdateOperationExtensions.kt
@@ -15,6 +15,8 @@
*/
package org.springframework.data.cassandra.core
+import kotlinx.coroutines.reactive.awaitSingle
+import org.springframework.data.cassandra.core.query.Update
import kotlin.reflect.KClass
/**
@@ -35,3 +37,11 @@ fun ReactiveUpdateOperation.update(entityClass: KClass): ReactiveUp
*/
inline fun ReactiveUpdateOperation.update(): ReactiveUpdateOperation.ReactiveUpdate =
update(T::class.java)
+
+/**
+ * Coroutines variant of [ReactiveUpdateOperation.TerminatingUpdate.apply].
+ *
+ * @author MarkPaluch
+ * @since 2.2
+ */
+suspend fun ReactiveUpdateOperation.TerminatingUpdate.applyAndAwait(update: Update): WriteResult = apply(update).awaitSingle()
diff --git a/spring-data-cassandra/src/test/kotlin/org/springframework/data/cassandra/core/ReactiveDeleteOperationExtensionsUnitTests.kt b/spring-data-cassandra/src/test/kotlin/org/springframework/data/cassandra/core/ReactiveDeleteOperationExtensionsUnitTests.kt
index b4c710cb7..0bbc05a40 100644
--- a/spring-data-cassandra/src/test/kotlin/org/springframework/data/cassandra/core/ReactiveDeleteOperationExtensionsUnitTests.kt
+++ b/spring-data-cassandra/src/test/kotlin/org/springframework/data/cassandra/core/ReactiveDeleteOperationExtensionsUnitTests.kt
@@ -15,10 +15,14 @@
*/
package org.springframework.data.cassandra.core
+import io.mockk.every
import io.mockk.mockk
import io.mockk.verify
+import kotlinx.coroutines.runBlocking
+import org.assertj.core.api.Assertions
import org.junit.Test
import org.springframework.data.cassandra.domain.Person
+import reactor.core.publisher.Mono
/**
* Unit tests for [ReactiveDeleteOperationExtensions].
@@ -42,4 +46,20 @@ class ReactiveDeleteOperationExtensionsUnitTests {
operations.delete()
verify { operations.delete(Person::class.java) }
}
+
+ @Test // DATACASS-632
+ fun allAndAwait() {
+
+ val delete = mockk()
+ val result = mockk()
+ every { delete.all() } returns Mono.just(result)
+
+ runBlocking {
+ Assertions.assertThat(delete.allAndAwait()).isEqualTo(result)
+ }
+
+ verify {
+ delete.all()
+ }
+ }
}
diff --git a/spring-data-cassandra/src/test/kotlin/org/springframework/data/cassandra/core/ReactiveInsertOperationExtensionsUnitTests.kt b/spring-data-cassandra/src/test/kotlin/org/springframework/data/cassandra/core/ReactiveInsertOperationExtensionsUnitTests.kt
index 62a700c8c..b9ac14ec6 100644
--- a/spring-data-cassandra/src/test/kotlin/org/springframework/data/cassandra/core/ReactiveInsertOperationExtensionsUnitTests.kt
+++ b/spring-data-cassandra/src/test/kotlin/org/springframework/data/cassandra/core/ReactiveInsertOperationExtensionsUnitTests.kt
@@ -15,11 +15,14 @@
*/
package org.springframework.data.cassandra.core
+import io.mockk.every
import io.mockk.mockk
import io.mockk.verify
-import org.junit.Ignore
+import kotlinx.coroutines.runBlocking
+import org.assertj.core.api.Assertions
import org.junit.Test
import org.springframework.data.cassandra.domain.Person
+import reactor.core.publisher.Mono
/**
* Unit tests for [ReactiveInsertOperationExtensions].
@@ -43,4 +46,20 @@ class ReactiveInsertOperationExtensionsUnitTests {
operations.insert()
verify { operations.insert(Person::class.java) }
}
+
+ @Test // DATACASS-632
+ fun oneAndAwait() {
+
+ val insert = mockk>()
+ val result = mockk>()
+ every { insert.one("foo") } returns Mono.just(result)
+
+ runBlocking {
+ Assertions.assertThat(insert.oneAndAwait("foo")).isEqualTo(result)
+ }
+
+ verify {
+ insert.one("foo")
+ }
+ }
}
diff --git a/spring-data-cassandra/src/test/kotlin/org/springframework/data/cassandra/core/ReactiveSelectOperationExtensionsUnitTests.kt b/spring-data-cassandra/src/test/kotlin/org/springframework/data/cassandra/core/ReactiveSelectOperationExtensionsUnitTests.kt
index d38a2fc2e..e60d4fe63 100644
--- a/spring-data-cassandra/src/test/kotlin/org/springframework/data/cassandra/core/ReactiveSelectOperationExtensionsUnitTests.kt
+++ b/spring-data-cassandra/src/test/kotlin/org/springframework/data/cassandra/core/ReactiveSelectOperationExtensionsUnitTests.kt
@@ -15,11 +15,15 @@
*/
package org.springframework.data.cassandra.core
+import io.mockk.every
import io.mockk.mockk
import io.mockk.verify
+import kotlinx.coroutines.runBlocking
+import org.assertj.core.api.Assertions
import org.junit.Test
import org.springframework.data.cassandra.domain.Person
import org.springframework.data.cassandra.domain.User
+import reactor.core.publisher.Mono
/**
* Unit tests for [ReactiveSelectOperationExtensions].
@@ -59,4 +63,64 @@ class ReactiveSelectOperationExtensionsUnitTests {
operationWithProjection.asType();
verify { operationWithProjection.`as`(User::class.java) }
}
+
+ @Test // DATACASS-632
+ fun terminatingFindAwaitOne() {
+
+ val find = mockk>()
+ every { find.one() } returns Mono.just("foo")
+
+ runBlocking {
+ Assertions.assertThat(find.awaitOne()).isEqualTo("foo")
+ }
+
+ verify {
+ find.one()
+ }
+ }
+
+ @Test // DATACASS-632
+ fun terminatingFindAwaitFirst() {
+
+ val find = mockk>()
+ every { find.first() } returns Mono.just("foo")
+
+ runBlocking {
+ Assertions.assertThat(find.awaitFirst()).isEqualTo("foo")
+ }
+
+ verify {
+ find.first()
+ }
+ }
+
+ @Test // DATACASS-632
+ fun terminatingFindAwaitCount() {
+
+ val find = mockk>()
+ every { find.count() } returns Mono.just(1)
+
+ runBlocking {
+ Assertions.assertThat(find.awaitCount()).isEqualTo(1)
+ }
+
+ verify {
+ find.count()
+ }
+ }
+
+ @Test // DATACASS-632
+ fun terminatingFindAwaitExists() {
+
+ val find = mockk>()
+ every { find.exists() } returns Mono.just(true)
+
+ runBlocking {
+ Assertions.assertThat(find.awaitExists()).isTrue()
+ }
+
+ verify {
+ find.exists()
+ }
+ }
}
diff --git a/spring-data-cassandra/src/test/kotlin/org/springframework/data/cassandra/core/ReactiveUpdateOperationExtensionsUnitTests.kt b/spring-data-cassandra/src/test/kotlin/org/springframework/data/cassandra/core/ReactiveUpdateOperationExtensionsUnitTests.kt
index 81c881e75..07652d3d7 100644
--- a/spring-data-cassandra/src/test/kotlin/org/springframework/data/cassandra/core/ReactiveUpdateOperationExtensionsUnitTests.kt
+++ b/spring-data-cassandra/src/test/kotlin/org/springframework/data/cassandra/core/ReactiveUpdateOperationExtensionsUnitTests.kt
@@ -15,10 +15,15 @@
*/
package org.springframework.data.cassandra.core
+import io.mockk.every
import io.mockk.mockk
import io.mockk.verify
+import kotlinx.coroutines.runBlocking
+import org.assertj.core.api.Assertions
import org.junit.Test
+import org.springframework.data.cassandra.core.query.Update
import org.springframework.data.cassandra.domain.Person
+import reactor.core.publisher.Mono
/**
* Unit tests for [ReactiveUpdateOperationExtensions].
@@ -42,4 +47,21 @@ class ReactiveUpdateOperationExtensionsUnitTests {
operations.update()
verify { operations.update(Person::class.java) }
}
+
+ @Test // DATACASS-632
+ fun applyAndAwait() {
+
+ val update = mockk()
+ val result = mockk()
+ val updateObj = mockk();
+ every { update.apply(updateObj) } returns Mono.just(result)
+
+ runBlocking {
+ Assertions.assertThat(update.applyAndAwait(updateObj)).isEqualTo(result)
+ }
+
+ verify {
+ update.apply(updateObj)
+ }
+ }
}
diff --git a/src/main/asciidoc/new-features.adoc b/src/main/asciidoc/new-features.adoc
index 50fa2cbbc..ee85997f3 100644
--- a/src/main/asciidoc/new-features.adoc
+++ b/src/main/asciidoc/new-features.adoc
@@ -7,6 +7,7 @@ This chapter summarizes changes and new features for each release.
== What's new in Spring Data for Apache Cassandra 2.2
* Read-only properties annotated with `@ReadOnlyProperty` to exclude non-writable properties from entity-bound `INSERT` and `UPDATE` operations.
* Support for derived `Between` queries.
+* Kotlin Coroutine extensions for `ReactiveFluentCassandraOperations`.
[[new-features.2-1-0]]
== What's new in Spring Data for Apache Cassandra 2.1