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]

Reply via email to