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]
