This is an automated email from the ASF dual-hosted git repository.

quantranhong1999 pushed a commit to branch 3.9.x
in repository https://gitbox.apache.org/repos/asf/james-project.git

commit 67701826a7f19d5dbff6c07b848787d4a7a09f2c
Author: Benoit TELLIER <[email protected]>
AuthorDate: Sun Aug 16 16:22:29 2026 +0700

    [FIX] Handle missing Task history in Cassandra implem like PG
---
 .../eventstore/cassandra/CassandraEventStore.scala    |  1 +
 .../cassandra/CassandraEventStoreExtension.scala      | 12 +++++++++---
 .../cassandra/CassandraEventStoreTest.scala           | 19 +++++++++++++++++++
 3 files changed, 29 insertions(+), 3 deletions(-)

diff --git 
a/event-sourcing/event-store-cassandra/src/main/scala/org/apache/james/eventsourcing/eventstore/cassandra/CassandraEventStore.scala
 
b/event-sourcing/event-store-cassandra/src/main/scala/org/apache/james/eventsourcing/eventstore/cassandra/CassandraEventStore.scala
index 9266a7e31a..2eef0a9816 100644
--- 
a/event-sourcing/event-store-cassandra/src/main/scala/org/apache/james/eventsourcing/eventstore/cassandra/CassandraEventStore.scala
+++ 
b/event-sourcing/event-store-cassandra/src/main/scala/org/apache/james/eventsourcing/eventstore/cassandra/CassandraEventStore.scala
@@ -49,6 +49,7 @@ class CassandraEventStore @Inject() (eventStoreDao: 
EventStoreDao) extends Event
   override def getEventsOfAggregate(aggregateId: AggregateId): SMono[History] =
     eventStoreDao.getSnapshot(aggregateId)
       .flatMap(snapshotId => eventStoreDao.getEventsOfAggregate(aggregateId, 
snapshotId))
+      .filter(history => history.getEvents.nonEmpty)
       .switchIfEmpty(eventStoreDao.getEventsOfAggregate(aggregateId))
 
   override def remove(aggregateId: AggregateId): Publisher[Void] = 
eventStoreDao.delete(aggregateId)
diff --git 
a/event-sourcing/event-store-cassandra/src/test/scala/org/apache/james/eventsourcing/eventstore/cassandra/CassandraEventStoreExtension.scala
 
b/event-sourcing/event-store-cassandra/src/test/scala/org/apache/james/eventsourcing/eventstore/cassandra/CassandraEventStoreExtension.scala
index d16dae4ee2..1662f22fab 100644
--- 
a/event-sourcing/event-store-cassandra/src/test/scala/org/apache/james/eventsourcing/eventstore/cassandra/CassandraEventStoreExtension.scala
+++ 
b/event-sourcing/event-store-cassandra/src/test/scala/org/apache/james/eventsourcing/eventstore/cassandra/CassandraEventStoreExtension.scala
@@ -18,6 +18,7 @@
  ****************************************************************/
 package org.apache.james.eventsourcing.eventstore.cassandra
 
+import com.datastax.oss.driver.api.core.CqlSession
 import org.apache.james.backends.cassandra.CassandraClusterExtension
 import org.apache.james.eventsourcing.eventstore.{EventStore, 
JsonEventSerializer}
 import org.junit.jupiter.api.extension.{AfterAllCallback, AfterEachCallback, 
BeforeAllCallback, BeforeEachCallback, ExtensionContext, ParameterContext, 
ParameterResolutionException, ParameterResolver}
@@ -40,9 +41,14 @@ class CassandraEventStoreExtension(var cassandra: 
CassandraClusterExtension, val
 
   @throws[ParameterResolutionException]
   override def supportsParameter(parameterContext: ParameterContext, 
extensionContext: ExtensionContext): Boolean =
-    parameterContext.getParameter.getType eq classOf[EventStore]
+    (parameterContext.getParameter.getType eq classOf[EventStore]) ||
+      (parameterContext.getParameter.getType eq classOf[CqlSession])
 
   @throws[ParameterResolutionException]
-  override def resolveParameter(parameterContext: ParameterContext, 
extensionContext: ExtensionContext): CassandraEventStore =
-    new CassandraEventStore(eventStoreDao.get)
+  override def resolveParameter(parameterContext: ParameterContext, 
extensionContext: ExtensionContext): AnyRef =
+    if (parameterContext.getParameter.getType eq classOf[CqlSession]) {
+      cassandra.getCassandraCluster.getConf
+    } else {
+      new CassandraEventStore(eventStoreDao.get)
+    }
 }
\ No newline at end of file
diff --git 
a/event-sourcing/event-store-cassandra/src/test/scala/org/apache/james/eventsourcing/eventstore/cassandra/CassandraEventStoreTest.scala
 
b/event-sourcing/event-store-cassandra/src/test/scala/org/apache/james/eventsourcing/eventstore/cassandra/CassandraEventStoreTest.scala
index 718f37c842..c55b42f824 100644
--- 
a/event-sourcing/event-store-cassandra/src/test/scala/org/apache/james/eventsourcing/eventstore/cassandra/CassandraEventStoreTest.scala
+++ 
b/event-sourcing/event-store-cassandra/src/test/scala/org/apache/james/eventsourcing/eventstore/cassandra/CassandraEventStoreTest.scala
@@ -18,6 +18,8 @@
  ****************************************************************/
 package org.apache.james.eventsourcing.eventstore.cassandra
 
+import com.datastax.oss.driver.api.core.CqlSession
+import com.datastax.oss.driver.api.querybuilder.QueryBuilder.{literal, update}
 import org.apache.james.eventsourcing.eventstore.dto.SnapshotEvent
 import org.apache.james.eventsourcing.{EventId, TestEvent}
 import org.apache.james.eventsourcing.eventstore.{EventStore, 
EventStoreContract, History}
@@ -55,4 +57,21 @@ class CassandraEventStoreTest extends EventStoreContract {
     
assertThat(SMono(testee.getEventsOfAggregate(EventStoreContract.AGGREGATE_1)).block())
       .isEqualTo(History.of(event3))
   }
+
+  @Test
+  def 
getEventsOfAggregateShouldFallBackToFullHistoryWhenSnapshotPointsToMissingEvents(testee:
 EventStore, session: CqlSession) : Unit = {
+    val event1 = TestEvent(EventId.first, EventStoreContract.AGGREGATE_1, 
"first")
+
+    SMono(testee.append(event1)).block()
+
+    // The snapshot is a static column written by a blind update, outside of 
the LWT appending the events:
+    // it can thus end up referencing events that are no longer part of the 
partition.
+    session.execute(update(CassandraEventStoreTable.EVENTS_TABLE)
+      .setColumn(CassandraEventStoreTable.SNAPSHOT, literal(5))
+      
.whereColumn(CassandraEventStoreTable.AGGREGATE_ID).isEqualTo(literal(EventStoreContract.AGGREGATE_1.asAggregateKey))
+      .build())
+
+    
assertThat(SMono(testee.getEventsOfAggregate(EventStoreContract.AGGREGATE_1)).block())
+      .isEqualTo(History.of(event1))
+  }
 }
\ No newline at end of file


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to