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 e5ac3864fb77f829124bd26e0862c82cf5c5ddff
Author: Benoit TELLIER <[email protected]>
AuthorDate: Sun Aug 16 16:25:34 2026 +0700

    [FIX] Cassandra event store: include the snapshot in the event batch
---
 .../eventstore/cassandra/EventStoreDao.scala             | 16 ++++++----------
 1 file changed, 6 insertions(+), 10 deletions(-)

diff --git 
a/event-sourcing/event-store-cassandra/src/main/scala/org/apache/james/eventsourcing/eventstore/cassandra/EventStoreDao.scala
 
b/event-sourcing/event-store-cassandra/src/main/scala/org/apache/james/eventsourcing/eventstore/cassandra/EventStoreDao.scala
index b2d7f11f16..63f013d925 100644
--- 
a/event-sourcing/event-store-cassandra/src/main/scala/org/apache/james/eventsourcing/eventstore/cassandra/EventStoreDao.scala
+++ 
b/event-sourcing/event-store-cassandra/src/main/scala/org/apache/james/eventsourcing/eventstore/cassandra/EventStoreDao.scala
@@ -89,20 +89,16 @@ class EventStoreDao @Inject() (val session: CqlSession,
       .build())
 
   private[cassandra] def appendAll(events: Iterable[Event], lastSnapShot: 
Option[EventId]): SMono[Boolean] =
-    SMono(cassandraAsyncExecutor.executeReturnApplied(appendQuery(events))
+    SMono(cassandraAsyncExecutor.executeReturnApplied(appendQuery(events, 
lastSnapShot))
       .map(_.booleanValue()))
-      .flatMap((success: Boolean) => lastSnapShot
-        .filter(_ => success)
-        .map(id => 
SMono(cassandraAsyncExecutor.executeVoid(insertSnapshot(events.head.getAggregateId,
 id))))
-        .getOrElse(SMono.empty)
-        .`then`(SMono.just(success)))
-
-  private def appendQuery(events: Iterable[Event]): Statement[_] =
-    if (events.size == 1)
+
+  private def appendQuery(events: Iterable[Event], lastSnapShot: 
Option[EventId]): Statement[_] =
+    if (events.size == 1 && lastSnapShot.isEmpty)
       insertEvent(events.head)
     else {
       val batch: BatchStatementBuilder = new 
BatchStatementBuilder(BatchType.LOGGED)
       events.foreach((event: Event) => batch.addStatement(insertEvent(event)))
+      lastSnapShot.foreach(snapshotId => 
batch.addStatement(insertSnapshot(events.head.getAggregateId, snapshotId)))
       batch.build()
     }
 
@@ -150,4 +146,4 @@ class EventStoreDao @Inject() (val session: CqlSession,
 
   private def toEvent(row: Row): SMono[Event] = SMono.fromCallable(() => 
jsonEventSerializer.deserialize(row.get(0, TypeCodecs.TEXT)))
     .subscribeOn(ReactorUtils.BLOCKING_CALL_WRAPPER)
-}
\ No newline at end of file
+}


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

Reply via email to