This is an automated email from the ASF dual-hosted git repository.
pjfanning pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/pekko-projection.git
The following commit(s) were added to refs/heads/main by this push:
new 269641d4 cassandra integration test issues (#593)
269641d4 is described below
commit 269641d41176719b476a4e8779d5c8d763668352
Author: PJ Fanning <[email protected]>
AuthorDate: Wed Jul 29 08:42:08 2026 +0100
cassandra integration test issues (#593)
* cassandra integration test issues
* Update ContainerSessionProvider.scala
* Update ContainerSessionProvider.scala
* Update CassandraProjectionTest.java
---
.../cassandra/CassandraProjectionTest.java | 49 ++++++++++++----------
.../cassandra/ContainerSessionProvider.scala | 8 +++-
2 files changed, 35 insertions(+), 22 deletions(-)
diff --git
a/cassandra-test/src/test/java/org/apache/pekko/projection/cassandra/CassandraProjectionTest.java
b/cassandra-test/src/test/java/org/apache/pekko/projection/cassandra/CassandraProjectionTest.java
index ad613862..25a7a462 100644
---
a/cassandra-test/src/test/java/org/apache/pekko/projection/cassandra/CassandraProjectionTest.java
+++
b/cassandra-test/src/test/java/org/apache/pekko/projection/cassandra/CassandraProjectionTest.java
@@ -14,7 +14,6 @@
package org.apache.pekko.projection.cassandra;
import static org.junit.jupiter.api.Assertions.assertEquals;
-import static org.junit.jupiter.api.Assertions.fail;
import java.time.Duration;
import java.util.List;
@@ -48,12 +47,14 @@ import
org.apache.pekko.projection.testkit.javadsl.TestSourceProvider;
import org.apache.pekko.stream.connectors.cassandra.javadsl.CassandraSession;
import
org.apache.pekko.stream.connectors.cassandra.javadsl.CassandraSessionRegistry;
import org.apache.pekko.stream.javadsl.Source;
+import org.apache.pekko.stream.testkit.TestSubscriber;
import org.junit.jupiter.api.AfterAll;
import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
import scala.concurrent.Await;
import scala.jdk.javaapi.FutureConverters;
+import scala.util.Either;
@ExtendWith(LogCapturingExtension.class)
public class CassandraProjectionTest {
@@ -191,6 +192,16 @@ public class CassandraProjectionTest {
});
}
+ @SuppressWarnings({"unchecked", "deprecation"})
+ private Throwable eventuallyExpectError(TestSubscriber.Probe<?> sinkProbe) {
+ while (true) {
+ Either<Throwable, ?> result = (Either<Throwable, ?>)
sinkProbe.expectNextOrError();
+ if (result.isLeft()) {
+ return result.left().get();
+ }
+ }
+ }
+
private Handler<Envelope> concatHandler(StringBuffer str) {
return Handler.fromFunction(
envelope -> {
@@ -262,17 +273,15 @@ public class CassandraProjectionTest {
projectionId, sourceProvider(entityId), () ->
concatHandlerFail4(str))
.withSaveOffset(1, Duration.ZERO);
- try {
- projectionTestKit.run(
- projection,
- () -> {
- assertEquals("abc|def|ghi|", str.toString());
- });
- fail("Expected exception");
- } catch (RuntimeException e) {
- assertEquals("fail on 4", e.getMessage());
- }
+ projectionTestKit.runWithTestSink(
+ projection,
+ sinkProbe -> {
+ sinkProbe.request(1000);
+ Throwable error = eventuallyExpectError(sinkProbe);
+ assertEquals("fail on 4", error.getMessage());
+ });
+ assertEquals("abc|def|ghi|", str.toString());
assertStoredOffset(projectionId, 3L);
// re-run projection without failing function
@@ -344,17 +353,15 @@ public class CassandraProjectionTest {
CassandraProjection.atMostOnce(
projectionId, sourceProvider(entityId), () ->
concatHandlerFail4(str));
- try {
- projectionTestKit.run(
- projection,
- () -> {
- assertEquals("abc|def|ghi|", str.toString());
- });
- fail("Expected exception");
- } catch (RuntimeException e) {
- assertEquals("fail on 4", e.getMessage());
- }
+ projectionTestKit.runWithTestSink(
+ projection,
+ sinkProbe -> {
+ sinkProbe.request(1000);
+ Throwable error = eventuallyExpectError(sinkProbe);
+ assertEquals("fail on 4", error.getMessage());
+ });
+ assertEquals("abc|def|ghi|", str.toString());
assertStoredOffset(projectionId, 4L);
// re-run projection without failing function
diff --git
a/cassandra-test/src/test/scala/org/apache/pekko/projection/cassandra/ContainerSessionProvider.scala
b/cassandra-test/src/test/scala/org/apache/pekko/projection/cassandra/ContainerSessionProvider.scala
index 2f63193c..847f00c8 100644
---
a/cassandra-test/src/test/scala/org/apache/pekko/projection/cassandra/ContainerSessionProvider.scala
+++
b/cassandra-test/src/test/scala/org/apache/pekko/projection/cassandra/ContainerSessionProvider.scala
@@ -19,23 +19,29 @@ import scala.concurrent.ExecutionContext
import scala.concurrent.Future
import scala.util.Try
+import org.apache.pekko.actor.ActorSystem
import org.apache.pekko.stream.connectors.cassandra.CqlSessionProvider
import com.datastax.oss.driver.api.core.CqlSession
+import
com.datastax.oss.driver.internal.core.config.typesafe.DefaultDriverConfigLoader
import com.datastax.oss.driver.internal.core.metadata.DefaultEndPoint
+import com.typesafe.config.Config
import org.testcontainers.cassandra.CassandraContainer
import org.testcontainers.utility.DockerImageName
/**
* Use testcontainers to lazily provide a single CqlSession for all Cassandra
tests
*/
-final class ContainerSessionProvider extends CqlSessionProvider {
+final class ContainerSessionProvider(system: ActorSystem, sessionConfig:
Config) extends CqlSessionProvider {
import ContainerSessionProvider._
+ private val driverConfig = CqlSessionProvider.driverConfig(system,
sessionConfig)
+
override def connect()(implicit ec: ExecutionContext): Future[CqlSession] =
started.map { _ =>
CqlSession.builder
.addContactEndPoint(new DefaultEndPoint(InetSocketAddress
.createUnresolved(container.getHost,
container.getFirstMappedPort.intValue())))
.withLocalDatacenter("datacenter1")
+ .withConfigLoader(new DefaultDriverConfigLoader(() => driverConfig))
.build()
}
}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]