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]
