From 6269042c75b3f1dbc6933827e24b90af3da2cd1d Mon Sep 17 00:00:00 2001 From: Olger Warnier Date: Sat, 5 Aug 2023 09:40:52 +0200 Subject: [PATCH 1/5] module akka to pekko --- .../src/main/resources/reference.conf | 0 .../akka/http/AkkaHttpParameterOffsetConverters.scala | 0 .../src/test/resources/logback-test.xml | 0 build.sbt | 9 ++++----- version.sbt | 2 +- 5 files changed, 5 insertions(+), 6 deletions(-) rename {bounded-akka-http => bounded-pekko-http}/src/main/resources/reference.conf (100%) rename {bounded-akka-http => bounded-pekko-http}/src/main/scala/io/cafienne/bounded/akka/http/AkkaHttpParameterOffsetConverters.scala (100%) rename {bounded-akka-http => bounded-pekko-http}/src/test/resources/logback-test.xml (100%) diff --git a/bounded-akka-http/src/main/resources/reference.conf b/bounded-pekko-http/src/main/resources/reference.conf similarity index 100% rename from bounded-akka-http/src/main/resources/reference.conf rename to bounded-pekko-http/src/main/resources/reference.conf diff --git a/bounded-akka-http/src/main/scala/io/cafienne/bounded/akka/http/AkkaHttpParameterOffsetConverters.scala b/bounded-pekko-http/src/main/scala/io/cafienne/bounded/akka/http/AkkaHttpParameterOffsetConverters.scala similarity index 100% rename from bounded-akka-http/src/main/scala/io/cafienne/bounded/akka/http/AkkaHttpParameterOffsetConverters.scala rename to bounded-pekko-http/src/main/scala/io/cafienne/bounded/akka/http/AkkaHttpParameterOffsetConverters.scala diff --git a/bounded-akka-http/src/test/resources/logback-test.xml b/bounded-pekko-http/src/test/resources/logback-test.xml similarity index 100% rename from bounded-akka-http/src/test/resources/logback-test.xml rename to bounded-pekko-http/src/test/resources/logback-test.xml diff --git a/build.sbt b/build.sbt index 1729776..ee017ba 100644 --- a/build.sbt +++ b/build.sbt @@ -1,12 +1,11 @@ lazy val basicSettings = { val scala213 = "2.13.11" - val scala212 = "2.12.13" - val supportedScalaVersions = List(scala213, scala212) + val supportedScalaVersions = List(scala213) Seq( organization := "io.cafienne.bounded", - description := "Scala and Akka based Domain Driven Design Framework", + description := "Scala and Pekko based Domain Driven Design Framework", scalaVersion := scala213, //crossScalaVersions := supportedScalaVersions, //releaseCrossBuild := false, @@ -75,12 +74,12 @@ val boundedCore = (project in file("bounded-core")) name := "bounded-core", libraryDependencies ++= Dependencies.baseDeps ++ Dependencies.persistanceLmdbDBDeps ++ Dependencies.persistenceCassandraDeps ++ Dependencies.persistenceJdbcDeps ++ Dependencies.testDeps) -val boundedAkkaHttp = (project in file("bounded-akka-http")) +val boundedAkkaHttp = (project in file("bounded-pekko-http")) .dependsOn(boundedCore) .enablePlugins(ReleasePlugin, AutomateHeaderPlugin) .settings(basicSettings: _*) .settings( - name := "bounded-akka-http", + name := "bounded-pekko-http", libraryDependencies ++= Dependencies.akkaHttpDeps) val boundedTest = (project in file("bounded-test")) diff --git a/version.sbt b/version.sbt index cba00a0..40bc6a2 100644 --- a/version.sbt +++ b/version.sbt @@ -1,2 +1,2 @@ -ThisBuild / version := "0.3.8" +ThisBuild / version := "0.4.0" From cf9dbb711e79344ee7811561cf685db56d0ce27f Mon Sep 17 00:00:00 2001 From: Olger Warnier Date: Tue, 26 Nov 2024 19:14:05 +0100 Subject: [PATCH 2/5] Move to pekko --- .../aggregate/AggregateRootActor.scala | 6 +- .../bounded/aggregate/CommandGateway.scala | 16 +- .../CommandValidationException.scala | 2 +- .../bounded/aggregate/CommandValidator.scala | 2 +- .../cafienne/bounded/aggregate/Protocol.scala | 4 +- .../typed/DefaultTypedCommandGateway.scala | 10 +- .../typed/TypedAggregateRootManager.scala | 6 +- .../bounded/akka/ActorSystemProvider.scala | 4 +- .../persistence/ReadJournalProvider.scala | 53 ++- .../persistence/leveldb/SharedJournal.scala | 6 +- .../cafienne/bounded/config/Configured.scala | 2 +- .../AbstractEventMaterializer.scala | 14 +- .../AbstractReplayableEventMaterializer.scala | 10 +- .../EventMaterializerExecutionContext.scala | 2 +- .../EventMaterializers.scala | 4 +- .../eventmaterializers/EventProcessed.scala | 4 +- .../eventmaterializers/ExceptionWriter.scala | 2 +- .../MaterializerEventFilter.scala | 2 +- .../OffsetStoreProvider.scala | 6 +- .../ReadJournalOffsetStore.scala | 11 +- .../ResumableReplayable.scala | 6 +- .../SagaEventMaterializer.scala | 16 +- .../offsetstores/CassandraOffsetStore.scala | 6 +- .../InMemoryBasedOffsetStore.scala | 6 +- .../offsetstores/JdbcOffsetStore.scala | 6 +- .../offsetstores/LmdbOffsetStore.scala | 4 +- .../offsetstores/OffsetStore.scala | 6 +- .../src/test/resources/logback-test.xml | 4 +- .../src/test/resources/reference.conf | 28 +- .../io/cafienne/bounded/SampleProtocol.scala | 2 +- .../aggregate/ClassicCommandGatewaySpec.scala | 32 +- .../aggregate/ClassicSimpleAggregate.scala | 4 +- .../aggregate/TypedAnotherAggregate.scala | 14 +- .../aggregate/TypedClusteredSpec.scala | 4 +- .../aggregate/TypedCommandGatewaySpec.scala | 10 +- .../aggregate/TypedSimpleAggregate.scala | 14 +- ...terializerWithRuntimeEventFilterSpec.scala | 352 +++++++++--------- .../CreateEventsInStoreActor.scala | 64 ++-- .../eventmaterializers/SpecConfig.scala | 144 ++++--- .../TestTaggingEventAdapter.scala | 34 +- .../AkkaHttpParameterOffsetConverters.scala | 6 +- .../src/test/resources/logback-test.xml | 2 +- .../test/ScalatestTypedActorHttpRoute.scala | 12 +- .../bounded/test/ClearStorageAfterEach.scala | 39 +- .../test/CreateEventsInStoreActor.scala | 6 +- .../bounded/test/LoggingTestProbe.scala | 6 +- .../bounded/test/StopSystemAfterAll.scala | 4 +- .../bounded/test/TestableAggregateRoot.scala | 30 +- .../bounded/test/TestableProjection.scala | 243 ++++++------ .../cafienne/bounded/test/TestableSaga.scala | 2 +- .../test/typed/TestTypedCommandGateway.scala | 10 +- .../test/typed/TestableAggregateRoot.scala | 44 ++- .../src/test/resources/logback-test.xml | 32 ++ .../src/test/resources/reference.conf | 20 + .../bounded/test/DomainProtocol.scala | 2 +- .../io/cafienne/bounded/test/SpecConfig.scala | 14 +- .../bounded/test/TestAggregateRoot.scala | 4 +- .../test/TestableAggregateRootSpec.scala | 14 +- .../bounded/test/TestableProjectionSpec.scala | 248 ++++++------ .../typed/TestableAggregateRootSpec.scala | 62 +-- .../test/typed/TypedSimpleAggregate.scala | 27 +- build.sbt | 32 +- project/Dependencies.scala | 70 ++-- project/build.properties | 2 +- project/plugins.sbt | 14 +- version.sbt | 2 +- 66 files changed, 961 insertions(+), 908 deletions(-) create mode 100644 bounded-test/src/test/resources/logback-test.xml create mode 100644 bounded-test/src/test/resources/reference.conf diff --git a/bounded-core/src/main/scala/io/cafienne/bounded/aggregate/AggregateRootActor.scala b/bounded-core/src/main/scala/io/cafienne/bounded/aggregate/AggregateRootActor.scala index 5189c44..d5e5ede 100644 --- a/bounded-core/src/main/scala/io/cafienne/bounded/aggregate/AggregateRootActor.scala +++ b/bounded-core/src/main/scala/io/cafienne/bounded/aggregate/AggregateRootActor.scala @@ -1,11 +1,11 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ package io.cafienne.bounded.aggregate -import akka.actor._ -import akka.persistence.{PersistentActor, RecoveryCompleted, SnapshotOffer} +import org.apache.pekko.actor._ +import org.apache.pekko.persistence.{PersistentActor, RecoveryCompleted, SnapshotOffer} import io.cafienne.bounded.aggregate.AggregateRootActor.GetState trait AggregateRootCreator { diff --git a/bounded-core/src/main/scala/io/cafienne/bounded/aggregate/CommandGateway.scala b/bounded-core/src/main/scala/io/cafienne/bounded/aggregate/CommandGateway.scala index ae522bf..30cf562 100644 --- a/bounded-core/src/main/scala/io/cafienne/bounded/aggregate/CommandGateway.scala +++ b/bounded-core/src/main/scala/io/cafienne/bounded/aggregate/CommandGateway.scala @@ -1,13 +1,13 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ package io.cafienne.bounded.aggregate -import akka.actor.{Actor, ActorContext, ActorRef, ActorSystem, Props, Terminated, TimerScheduler, Timers} -import akka.event.{Logging, LoggingAdapter} -import akka.routing.{NoRoutee, Routee, RoutingLogic} -import akka.util.Timeout +import org.apache.pekko.actor.{Actor, ActorContext, ActorRef, ActorSystem, Props, Terminated, TimerScheduler, Timers} +import org.apache.pekko.event.{Logging, LoggingAdapter} +import org.apache.pekko.routing.{NoRoutee, Routee, RoutingLogic} +import org.apache.pekko.util.Timeout import scala.collection.immutable import scala.concurrent.duration._ import scala.concurrent.{ExecutionContext, Future} @@ -21,7 +21,7 @@ class DefaultCommandGateway[A <: AggregateRootCreator](system: ActorSystem, aggr implicit timeout: Timeout, ec: ExecutionContext ) extends CommandGateway { - import akka.pattern.ask + import org.apache.pekko.pattern.ask implicit val actorSystem: ActorSystem = system val logger: LoggingAdapter = Logging(system, getClass) @@ -68,7 +68,7 @@ class RouterCommandGateway[A <: AggregateRootCreator]( implicit timeout: Timeout, ec: ExecutionContext ) extends CommandGateway { - import akka.pattern.ask + import org.apache.pekko.pattern.ask implicit val actorSystem: ActorSystem = system val logger: LoggingAdapter = Logging(system, getClass) @@ -90,7 +90,7 @@ class RouterCommandGateway[A <: AggregateRootCreator]( } -import akka.routing.{ActorRefRoutee, Router} +import org.apache.pekko.routing.{ActorRefRoutee, Router} private class AggregateGatewayRouter[A <: AggregateRootCreator](aggregateRootCreator: A, idleTime: FiniteDuration) extends Actor diff --git a/bounded-core/src/main/scala/io/cafienne/bounded/aggregate/CommandValidationException.scala b/bounded-core/src/main/scala/io/cafienne/bounded/aggregate/CommandValidationException.scala index 8565524..e7829fa 100644 --- a/bounded-core/src/main/scala/io/cafienne/bounded/aggregate/CommandValidationException.scala +++ b/bounded-core/src/main/scala/io/cafienne/bounded/aggregate/CommandValidationException.scala @@ -1,5 +1,5 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ package io.cafienne.bounded.aggregate diff --git a/bounded-core/src/main/scala/io/cafienne/bounded/aggregate/CommandValidator.scala b/bounded-core/src/main/scala/io/cafienne/bounded/aggregate/CommandValidator.scala index f2b17c3..b75a0ee 100644 --- a/bounded-core/src/main/scala/io/cafienne/bounded/aggregate/CommandValidator.scala +++ b/bounded-core/src/main/scala/io/cafienne/bounded/aggregate/CommandValidator.scala @@ -1,5 +1,5 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ package io.cafienne.bounded.aggregate diff --git a/bounded-core/src/main/scala/io/cafienne/bounded/aggregate/Protocol.scala b/bounded-core/src/main/scala/io/cafienne/bounded/aggregate/Protocol.scala index a0a4406..e55b56e 100644 --- a/bounded-core/src/main/scala/io/cafienne/bounded/aggregate/Protocol.scala +++ b/bounded-core/src/main/scala/io/cafienne/bounded/aggregate/Protocol.scala @@ -1,10 +1,10 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ package io.cafienne.bounded.aggregate -import akka.actor.typed.ActorRef +import org.apache.pekko.actor.typed.ActorRef import scala.collection.immutable.Seq trait DomainCommand { diff --git a/bounded-core/src/main/scala/io/cafienne/bounded/aggregate/typed/DefaultTypedCommandGateway.scala b/bounded-core/src/main/scala/io/cafienne/bounded/aggregate/typed/DefaultTypedCommandGateway.scala index 1eec664..cfec9fd 100644 --- a/bounded-core/src/main/scala/io/cafienne/bounded/aggregate/typed/DefaultTypedCommandGateway.scala +++ b/bounded-core/src/main/scala/io/cafienne/bounded/aggregate/typed/DefaultTypedCommandGateway.scala @@ -1,12 +1,12 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ package io.cafienne.bounded.aggregate.typed -import akka.actor.typed.{ActorRef, ActorSystem, Behavior, PostStop, Scheduler, Terminated} -import akka.actor.typed.scaladsl.{Behaviors, TimerScheduler} -import akka.util.Timeout +import org.apache.pekko.actor.typed.{ActorRef, ActorSystem, Behavior, PostStop, Scheduler, Terminated} +import org.apache.pekko.actor.typed.scaladsl.{Behaviors, TimerScheduler} +import org.apache.pekko.util.Timeout import io.cafienne.bounded.aggregate.{CommandValidator, DomainCommand, ValidateableCommand} import scala.concurrent.duration._ import scala.concurrent.{ExecutionContext, Future} @@ -42,7 +42,7 @@ class DefaultTypedCommandGateway[Cmd <: DomainCommand]( )(implicit timeout: Timeout, ec: ExecutionContext, scheduler: Scheduler) extends TypedCommandGateway[Cmd] { - import akka.actor.typed.scaladsl.AskPattern._ + import org.apache.pekko.actor.typed.scaladsl.AskPattern._ object CommandGatewayGuardian { diff --git a/bounded-core/src/main/scala/io/cafienne/bounded/aggregate/typed/TypedAggregateRootManager.scala b/bounded-core/src/main/scala/io/cafienne/bounded/aggregate/typed/TypedAggregateRootManager.scala index f40b500..98f5827 100644 --- a/bounded-core/src/main/scala/io/cafienne/bounded/aggregate/typed/TypedAggregateRootManager.scala +++ b/bounded-core/src/main/scala/io/cafienne/bounded/aggregate/typed/TypedAggregateRootManager.scala @@ -1,11 +1,11 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ package io.cafienne.bounded.aggregate.typed -import akka.actor.typed.Behavior -import akka.cluster.sharding.typed.scaladsl.EntityTypeKey +import org.apache.pekko.actor.typed.Behavior +import org.apache.pekko.cluster.sharding.typed.scaladsl.EntityTypeKey import io.cafienne.bounded.aggregate.DomainCommand trait TypedAggregateRootManager[T <: DomainCommand] { diff --git a/bounded-core/src/main/scala/io/cafienne/bounded/akka/ActorSystemProvider.scala b/bounded-core/src/main/scala/io/cafienne/bounded/akka/ActorSystemProvider.scala index b826ed7..dae3e9d 100644 --- a/bounded-core/src/main/scala/io/cafienne/bounded/akka/ActorSystemProvider.scala +++ b/bounded-core/src/main/scala/io/cafienne/bounded/akka/ActorSystemProvider.scala @@ -1,10 +1,10 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ package io.cafienne.bounded.akka -import akka.actor.ActorSystem +import org.apache.pekko.actor.ActorSystem trait ActorSystemProvider { implicit def system: ActorSystem diff --git a/bounded-core/src/main/scala/io/cafienne/bounded/akka/persistence/ReadJournalProvider.scala b/bounded-core/src/main/scala/io/cafienne/bounded/akka/persistence/ReadJournalProvider.scala index e09d1ca..4d58d43 100644 --- a/bounded-core/src/main/scala/io/cafienne/bounded/akka/persistence/ReadJournalProvider.scala +++ b/bounded-core/src/main/scala/io/cafienne/bounded/akka/persistence/ReadJournalProvider.scala @@ -1,28 +1,27 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ package io.cafienne.bounded.akka.persistence -import akka.actor.Props -import akka.persistence.cassandra.query.scaladsl.CassandraReadJournal -import akka.persistence.inmemory.query.scaladsl.InMemoryReadJournal -import akka.persistence.jdbc.query.scaladsl.JdbcReadJournal -import akka.persistence.journal.leveldb.{SharedLeveldbJournal, SharedLeveldbStore} -import akka.persistence.query.PersistenceQuery -import akka.persistence.query.journal.leveldb.scaladsl.LeveldbReadJournal -import akka.persistence.query.scaladsl._ +import org.apache.pekko.persistence.cassandra.query.scaladsl.CassandraReadJournal +import org.apache.pekko.persistence.query.PersistenceQuery +import org.apache.pekko.persistence.query.journal.leveldb.scaladsl.LeveldbReadJournal +import org.apache.pekko.persistence.query.scaladsl._ import io.cafienne.bounded.akka.ActorSystemProvider -import io.cafienne.bounded.akka.persistence.leveldb.SharedJournal +import org.apache.pekko.persistence.jdbc.query.scaladsl.JdbcReadJournal +import org.apache.pekko.persistence.r2dbc.query.scaladsl.R2dbcReadJournal +import org.apache.pekko.persistence.journal.inmem.InmemJournal +import org.apache.pekko.persistence.testkit.PersistenceTestKitPlugin +import org.apache.pekko.persistence.testkit.query.scaladsl.PersistenceTestKitReadJournal /** * Provides a readJournal that has the eventsByTag available that's used for * creation of the domain/query models of the system. */ trait ReadJournalProvider { systemProvider: ActorSystemProvider => - val configuredJournal = - system.settings.config.getString("akka.persistence.journal.plugin") + system.settings.config.getString("pekko.persistence.journal.plugin") def readJournal : ReadJournal with CurrentEventsByTagQuery with EventsByTagQuery with CurrentEventsByPersistenceIdQuery = { @@ -32,31 +31,14 @@ trait ReadJournalProvider { systemProvider: ActorSystemProvider => return PersistenceQuery(system) .readJournalFor[LeveldbReadJournal](LeveldbReadJournal.Identifier) } - if (configuredJournal.endsWith("leveldb-shared")) { - system.log.debug("configuring read journal for leveldb-shared") - - val sharedJournal = - system.actorOf(Props(new SharedLeveldbStore), SharedJournal.name) - SharedLeveldbJournal.setStore(sharedJournal, system) - - return PersistenceQuery(system) - .readJournalFor[LeveldbReadJournal](LeveldbReadJournal.Identifier) - } - //0.x series - if (configuredJournal.endsWith("cassandra-journal")) { + if (configuredJournal.endsWith("cassandra-journal") || configuredJournal.endsWith("cassandra.journal")) { system.log.debug("configuring read journal for cassandra") return PersistenceQuery(system) .readJournalFor[CassandraReadJournal](CassandraReadJournal.Identifier) } - // 1.x series - if (configuredJournal.endsWith("cassandra.journal")) { - system.log.debug("configuring read journal for cassandra") - return PersistenceQuery(system) - .readJournalFor[CassandraReadJournal](CassandraReadJournal.Identifier) - } - if (configuredJournal.endsWith("inmemory-journal")) { + if (configuredJournal.endsWith("r2dbc-journal")) { return PersistenceQuery(system) - .readJournalFor[InMemoryReadJournal](InMemoryReadJournal.Identifier) + .readJournalFor[R2dbcReadJournal](R2dbcReadJournal.Identifier) .asInstanceOf[ ReadJournal with CurrentPersistenceIdsQuery with CurrentEventsByPersistenceIdQuery with CurrentEventsByTagQuery with EventsByPersistenceIdQuery with EventsByTagQuery ] @@ -64,6 +46,13 @@ trait ReadJournalProvider { systemProvider: ActorSystemProvider => if (configuredJournal.endsWith("jdbc-journal")) { return PersistenceQuery(system) .readJournalFor[JdbcReadJournal](JdbcReadJournal.Identifier) +// .asInstanceOf[ +// ReadJournal with CurrentPersistenceIdsQuery with CurrentEventsByPersistenceIdQuery with CurrentEventsByTagQuery with EventsByPersistenceIdQuery with EventsByTagQuery +// ] + } + if (configuredJournal.endsWith("inmem")) { + return PersistenceQuery(system) + .readJournalFor("pekko.persistence.journal.inmem") .asInstanceOf[ ReadJournal with CurrentPersistenceIdsQuery with CurrentEventsByPersistenceIdQuery with CurrentEventsByTagQuery with EventsByPersistenceIdQuery with EventsByTagQuery ] diff --git a/bounded-core/src/main/scala/io/cafienne/bounded/akka/persistence/leveldb/SharedJournal.scala b/bounded-core/src/main/scala/io/cafienne/bounded/akka/persistence/leveldb/SharedJournal.scala index 22ac869..2520b0d 100644 --- a/bounded-core/src/main/scala/io/cafienne/bounded/akka/persistence/leveldb/SharedJournal.scala +++ b/bounded-core/src/main/scala/io/cafienne/bounded/akka/persistence/leveldb/SharedJournal.scala @@ -1,11 +1,11 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ package io.cafienne.bounded.akka.persistence.leveldb -import akka.actor.{ActorPath, RootActorPath, Address} -import akka.persistence.query.journal.leveldb.scaladsl.LeveldbReadJournal +import org.apache.pekko.actor.{ActorPath, RootActorPath, Address} +import org.apache.pekko.persistence.query.journal.leveldb.scaladsl.LeveldbReadJournal object SharedJournal { diff --git a/bounded-core/src/main/scala/io/cafienne/bounded/config/Configured.scala b/bounded-core/src/main/scala/io/cafienne/bounded/config/Configured.scala index dd8a85b..1641541 100644 --- a/bounded-core/src/main/scala/io/cafienne/bounded/config/Configured.scala +++ b/bounded-core/src/main/scala/io/cafienne/bounded/config/Configured.scala @@ -1,5 +1,5 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ package io.cafienne.bounded.config diff --git a/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/AbstractEventMaterializer.scala b/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/AbstractEventMaterializer.scala index 884007d..fad91cc 100644 --- a/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/AbstractEventMaterializer.scala +++ b/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/AbstractEventMaterializer.scala @@ -1,15 +1,15 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ package io.cafienne.bounded.eventmaterializers import java.util.UUID -import akka.{Done, NotUsed} -import akka.actor.ActorSystem -import akka.persistence.query.{EventEnvelope, Offset} -import akka.stream.{KillSwitches, UniqueKillSwitch} -import akka.stream.scaladsl.{Keep, Sink, Source} +import org.apache.pekko.{Done, NotUsed} +import org.apache.pekko.actor.ActorSystem +import org.apache.pekko.persistence.query.{EventEnvelope, Offset} +import org.apache.pekko.stream.{KillSwitches, UniqueKillSwitch} +import org.apache.pekko.stream.scaladsl.{Keep, Sink, Source} import io.cafienne.bounded.akka.ActorSystemProvider import io.cafienne.bounded.akka.persistence.ReadJournalProvider import io.cafienne.bounded.config.Configured @@ -18,7 +18,7 @@ import io.cafienne.bounded.aggregate.DomainEvent import io.cafienne.bounded.eventmaterializers.offsetstores.OffsetStore import scala.concurrent.{Await, ExecutionContextExecutor, Future} -import scala.concurrent.duration.* +import scala.concurrent.duration._ /** * Abstract class to be used to create eventlisteners. Use this class as a base for listening for diff --git a/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/AbstractReplayableEventMaterializer.scala b/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/AbstractReplayableEventMaterializer.scala index f16de65..0831454 100644 --- a/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/AbstractReplayableEventMaterializer.scala +++ b/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/AbstractReplayableEventMaterializer.scala @@ -1,13 +1,13 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ package io.cafienne.bounded.eventmaterializers -import akka.{Done, NotUsed} -import akka.actor.ActorSystem -import akka.persistence.query.{EventEnvelope, Offset} -import akka.stream.scaladsl.Source +import org.apache.pekko.{Done, NotUsed} +import org.apache.pekko.actor.ActorSystem +import org.apache.pekko.persistence.query.{EventEnvelope, Offset} +import org.apache.pekko.stream.scaladsl.Source import io.cafienne.bounded.aggregate.DomainEvent import scala.concurrent.Future diff --git a/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/EventMaterializerExecutionContext.scala b/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/EventMaterializerExecutionContext.scala index d77f9eb..b63ac43 100644 --- a/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/EventMaterializerExecutionContext.scala +++ b/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/EventMaterializerExecutionContext.scala @@ -1,5 +1,5 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ package io.cafienne.bounded.eventmaterializers diff --git a/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/EventMaterializers.scala b/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/EventMaterializers.scala index 6688648..5b2e981 100644 --- a/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/EventMaterializers.scala +++ b/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/EventMaterializers.scala @@ -1,10 +1,10 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ package io.cafienne.bounded.eventmaterializers -import akka.persistence.query.{Offset} +import org.apache.pekko.persistence.query.{Offset} import com.typesafe.scalalogging.Logger import org.slf4j.LoggerFactory diff --git a/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/EventProcessed.scala b/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/EventProcessed.scala index 71b2226..283e568 100644 --- a/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/EventProcessed.scala +++ b/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/EventProcessed.scala @@ -1,10 +1,10 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ package io.cafienne.bounded.eventmaterializers import java.util.UUID -import akka.persistence.query.Offset +import org.apache.pekko.persistence.query.Offset case class EventProcessed(materializerId: UUID, offset: Offset, persistenceId: String, sequenceNr: Long, evt: Any) diff --git a/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/ExceptionWriter.scala b/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/ExceptionWriter.scala index 7e6f9a6..16bc7e5 100644 --- a/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/ExceptionWriter.scala +++ b/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/ExceptionWriter.scala @@ -1,5 +1,5 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ package io.cafienne.bounded.eventmaterializers diff --git a/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/MaterializerEventFilter.scala b/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/MaterializerEventFilter.scala index 9456ca4..650de01 100644 --- a/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/MaterializerEventFilter.scala +++ b/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/MaterializerEventFilter.scala @@ -1,5 +1,5 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ package io.cafienne.bounded.eventmaterializers diff --git a/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/OffsetStoreProvider.scala b/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/OffsetStoreProvider.scala index 16d7504..a8ee883 100644 --- a/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/OffsetStoreProvider.scala +++ b/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/OffsetStoreProvider.scala @@ -1,11 +1,11 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ package io.cafienne.bounded.eventmaterializers -import akka.Done -import akka.persistence.query.Offset +import org.apache.pekko.Done +import org.apache.pekko.persistence.query.Offset import com.typesafe.config.Config import io.cafienne.bounded.akka.ActorSystemProvider import io.cafienne.bounded.eventmaterializers.offsetstores.{ diff --git a/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/ReadJournalOffsetStore.scala b/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/ReadJournalOffsetStore.scala index 9111d0f..061cbda 100644 --- a/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/ReadJournalOffsetStore.scala +++ b/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/ReadJournalOffsetStore.scala @@ -1,13 +1,14 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ package io.cafienne.bounded.eventmaterializers -import java.util.concurrent.TimeUnit +import org.apache.pekko.persistence.cassandra.query.scaladsl.CassandraReadJournal -import akka.persistence.cassandra.query.scaladsl.CassandraReadJournal -import akka.persistence.query.Offset +import java.util.concurrent.TimeUnit +import org.apache.pekko.persistence.cassandra.query.scaladsl.CassandraReadJournal +import org.apache.pekko.persistence.query.Offset import com.typesafe.config.{ConfigFactory, ConfigValueFactory} import io.cafienne.bounded.akka.ActorSystemProvider import io.cafienne.bounded.akka.persistence.ReadJournalProvider @@ -46,7 +47,7 @@ trait ReadJournalOffsetStore extends OffsetStore { new JdbcOffsetStore( DatabaseConfig.forConfig( config = system.settings.config, - path = system.settings.config.getString("akka.persistence.offset.jdbc.store") + path = system.settings.config.getString("akka.persistence.offset.r2dbc.store") ) ) } else { diff --git a/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/ResumableReplayable.scala b/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/ResumableReplayable.scala index 85289bd..67aae52 100644 --- a/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/ResumableReplayable.scala +++ b/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/ResumableReplayable.scala @@ -1,11 +1,11 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ package io.cafienne.bounded.eventmaterializers -import akka.Done -import akka.persistence.query.Offset +import org.apache.pekko.Done +import org.apache.pekko.persistence.query.Offset import scala.concurrent.Future diff --git a/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/SagaEventMaterializer.scala b/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/SagaEventMaterializer.scala index b585feb..e07ef77 100644 --- a/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/SagaEventMaterializer.scala +++ b/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/SagaEventMaterializer.scala @@ -1,21 +1,21 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ package io.cafienne.bounded.eventmaterializers import java.util.UUID -import akka.actor.ActorSystem -import akka.persistence.query.scaladsl.{ +import org.apache.pekko.actor.ActorSystem +import org.apache.pekko.persistence.query.scaladsl.{ CurrentEventsByPersistenceIdQuery, CurrentEventsByTagQuery, EventsByTagQuery, ReadJournal } -import akka.persistence.query.{EventEnvelope, Offset} -import akka.stream.scaladsl.{Keep, Merge, Sink, Source} -import akka.stream.{ActorMaterializer, KillSwitches, UniqueKillSwitch} -import akka.{Done, NotUsed} +import org.apache.pekko.persistence.query.{EventEnvelope, Offset} +import org.apache.pekko.stream.scaladsl.{Keep, Merge, Sink, Source} +import org.apache.pekko.stream.{ActorMaterializer, KillSwitches, UniqueKillSwitch} +import org.apache.pekko.{Done, NotUsed} import com.typesafe.scalalogging.Logger import io.cafienne.bounded.aggregate.DomainEvent import io.cafienne.bounded.akka.ActorSystemProvider @@ -23,7 +23,7 @@ import io.cafienne.bounded.akka.persistence.ReadJournalProvider import io.cafienne.bounded.config.Configured import io.cafienne.bounded.eventmaterializers.offsetstores.OffsetStore -import scala.concurrent.duration.* +import scala.concurrent.duration._ import scala.concurrent.{Await, ExecutionContextExecutor, Future} /** diff --git a/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/offsetstores/CassandraOffsetStore.scala b/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/offsetstores/CassandraOffsetStore.scala index a636cf5..84cd449 100644 --- a/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/offsetstores/CassandraOffsetStore.scala +++ b/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/offsetstores/CassandraOffsetStore.scala @@ -1,13 +1,13 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ package io.cafienne.bounded.eventmaterializers.offsetstores import java.util.UUID -import akka.persistence.cassandra.query.scaladsl.CassandraReadJournal -import akka.persistence.query.{Offset, Sequence, TimeBasedUUID} +import org.apache.pekko.persistence.cassandra.query.scaladsl.CassandraReadJournal +import org.apache.pekko.persistence.query.{Offset, Sequence, TimeBasedUUID} import io.cafienne.bounded.eventmaterializers.{EventMaterializerExecutionContext} import scala.concurrent.duration._ diff --git a/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/offsetstores/InMemoryBasedOffsetStore.scala b/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/offsetstores/InMemoryBasedOffsetStore.scala index 2da1e55..9529114 100644 --- a/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/offsetstores/InMemoryBasedOffsetStore.scala +++ b/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/offsetstores/InMemoryBasedOffsetStore.scala @@ -1,11 +1,11 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ package io.cafienne.bounded.eventmaterializers.offsetstores -import akka.Done -import akka.persistence.query.Offset +import org.apache.pekko.Done +import org.apache.pekko.persistence.query.Offset import scala.concurrent.Future diff --git a/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/offsetstores/JdbcOffsetStore.scala b/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/offsetstores/JdbcOffsetStore.scala index 53e0238..4cb6300 100644 --- a/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/offsetstores/JdbcOffsetStore.scala +++ b/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/offsetstores/JdbcOffsetStore.scala @@ -1,10 +1,10 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ package io.cafienne.bounded.eventmaterializers.offsetstores -import akka.Done +import org.apache.pekko.Done import slick.basic.DatabaseConfig import slick.jdbc.JdbcProfile import slick.lifted.ProvenShape @@ -14,7 +14,7 @@ import scala.concurrent.ExecutionContext import java.util.UUID import java.util.concurrent.Executors -import akka.persistence.query.{ +import org.apache.pekko.persistence.query.{ NoOffset => PersistenceNoOffset, Offset => PersistenceOffset, Sequence => PersistenceSequence, diff --git a/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/offsetstores/LmdbOffsetStore.scala b/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/offsetstores/LmdbOffsetStore.scala index 9dbdcbe..4170ef8 100644 --- a/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/offsetstores/LmdbOffsetStore.scala +++ b/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/offsetstores/LmdbOffsetStore.scala @@ -1,5 +1,5 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ package io.cafienne.bounded.eventmaterializers.offsetstores @@ -10,7 +10,7 @@ import java.nio.ByteBuffer.allocateDirect import java.nio.charset.StandardCharsets.UTF_8 import java.util.UUID -import akka.persistence.query.{Offset, Sequence, TimeBasedUUID} +import org.apache.pekko.persistence.query.{Offset, Sequence, TimeBasedUUID} import com.typesafe.config.Config import org.lmdbjava.{Dbi, DbiFlags, Env, EnvFlags} import org.slf4j.LoggerFactory diff --git a/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/offsetstores/OffsetStore.scala b/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/offsetstores/OffsetStore.scala index 584fbc3..ff05d5b 100644 --- a/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/offsetstores/OffsetStore.scala +++ b/bounded-core/src/main/scala/io/cafienne/bounded/eventmaterializers/offsetstores/OffsetStore.scala @@ -1,11 +1,11 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ package io.cafienne.bounded.eventmaterializers.offsetstores -import akka.Done -import akka.persistence.query.Offset +import org.apache.pekko.Done +import org.apache.pekko.persistence.query.Offset import scala.concurrent.Future diff --git a/bounded-core/src/test/resources/logback-test.xml b/bounded-core/src/test/resources/logback-test.xml index b75e8ba..5c0638f 100644 --- a/bounded-core/src/test/resources/logback-test.xml +++ b/bounded-core/src/test/resources/logback-test.xml @@ -15,14 +15,14 @@ the captured logging events are flushed to the appenders defined for the akka.actor.testkit.typed.internal.CapturingAppenderDelegate logger. --> - + - + diff --git a/bounded-core/src/test/resources/reference.conf b/bounded-core/src/test/resources/reference.conf index e9ccf85..2894f97 100644 --- a/bounded-core/src/test/resources/reference.conf +++ b/bounded-core/src/test/resources/reference.conf @@ -1,37 +1,17 @@ -akka { +pekko { loglevel = "DEBUG" stdout-loglevel = "DEBUG" - loggers = ["akka.testkit.TestEventListener"] + loggers = ["org.apache.pekko.testkit.TestEventListener"] actor { allow-java-serialization = on } -// actor { -// default-dispatcher { -// executor = "fork-join-executor" -// fork-join-executor { -// parallelism-min = 8 -// parallelism-factor = 2.0 -// parallelism-max = 8 -// } -// } -// serialize-creators = off -// serialize-messages = off -// serializers { -// serializer = "com.eduarte.participantplanner.persistence.ParticipantPlannerPersistersSerializer" -// } -// serialization-bindings { -// "stamina.Persistable" = serializer -// // enable below to check if all events have been serialized without java.io.Serializable -// "java.io.Serializable" = none -// } -// } persistence { publish-confirmations = on publish-plugin-commands = on journal { - plugin = "inmemory-journal" + plugin = "pekko.persistence.journal.inmem" } - snapshot-store.plugin = "inmemory-snapshot-store" + snapshot-store.plugin = "pekko.persistence.snapshot-store.local" } // test { // single-expect-default = 10s diff --git a/bounded-core/src/test/scala/io/cafienne/bounded/SampleProtocol.scala b/bounded-core/src/test/scala/io/cafienne/bounded/SampleProtocol.scala index cbd9193..e777561 100644 --- a/bounded-core/src/test/scala/io/cafienne/bounded/SampleProtocol.scala +++ b/bounded-core/src/test/scala/io/cafienne/bounded/SampleProtocol.scala @@ -1,5 +1,5 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ package io.cafienne.bounded diff --git a/bounded-core/src/test/scala/io/cafienne/bounded/aggregate/ClassicCommandGatewaySpec.scala b/bounded-core/src/test/scala/io/cafienne/bounded/aggregate/ClassicCommandGatewaySpec.scala index 29c4ac0..4af519c 100644 --- a/bounded-core/src/test/scala/io/cafienne/bounded/aggregate/ClassicCommandGatewaySpec.scala +++ b/bounded-core/src/test/scala/io/cafienne/bounded/aggregate/ClassicCommandGatewaySpec.scala @@ -1,12 +1,12 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ package io.cafienne.bounded.aggregate -import akka.actor.ActorSystem -import akka.testkit.{EventFilter, TestKit} -import akka.util.Timeout +import org.apache.pekko.actor.ActorSystem +import org.apache.pekko.testkit.{EventFilter, TestKit} +import org.apache.pekko.util.Timeout import com.typesafe.config.ConfigFactory import org.scalatest.BeforeAndAfterAll import org.scalatest.concurrent.ScalaFutures @@ -14,7 +14,7 @@ import org.scalatest.flatspec.AsyncFlatSpecLike import org.scalatest.matchers.should.Matchers import scala.concurrent.{Await, ExecutionContext, ExecutionContextExecutor, Future} -import scala.concurrent.duration.* +import scala.concurrent.duration._ import org.scalatest.time.{Millis, Seconds, Span} class ClassicCommandGatewaySpec @@ -22,19 +22,19 @@ class ClassicCommandGatewaySpec ActorSystem( "ClassicCommandGatewaySpec", ConfigFactory.parseString(s""" - akka.persistence.publish-plugin-commands = on - akka.persistence.journal.plugin = "akka.persistence.journal.inmem" - akka.persistence.journal.inmem.test-serialization = on - akka.actor.warn-about-java-serializer-usage = off + pekko.persistence.publish-plugin-commands = on + pekko.persistence.journal.plugin = "pekko.persistence.journal.inmem" + pekko.persistence.journal.inmem.test-serialization = on + pekko.actor.warn-about-java-serializer-usage = off # snapshot store plugin is NOT defined, things should still work - akka.persistence.snapshot-store.local.dir = "target/snapshots-${classOf[TypedCommandGatewaySpec].getName}/" + pekko.persistence.snapshot-store.local.dir = "target/snapshots-${classOf[TypedCommandGatewaySpec].getName}/" #PLEASE NOTE THAT CoordinatedShutdown needs to be disabled as below in order to run the test properly #SEE https://doc.akka.io/docs/akka/current/coordinated-shutdown.html at the end of the page - akka.coordinated-shutdown.terminate-actor-system = off - akka.coordinated-shutdown.run-by-actor-system-terminate = off - akka.coordinated-shutdown.run-by-jvm-shutdown-hook = off - akka.cluster.run-coordinated-shutdown-when-down = off - akka.loggers = ["akka.testkit.TestEventListener"] + pekko.coordinated-shutdown.terminate-actor-system = off + pekko.coordinated-shutdown.run-by-actor-system-terminate = off + pekko.coordinated-shutdown.run-by-jvm-shutdown-hook = off + pekko.cluster.run-coordinated-shutdown-when-down = off + pekko.loggers = ["org.apache.pekko.testkit.TestEventListener"] """) ) ) @@ -53,7 +53,7 @@ class ClassicCommandGatewaySpec val commandGateway: RouterCommandGateway[ClassicSimpleAggregateCreator] = new RouterCommandGateway(system, 2.seconds, classicSimpleCreator) - import akka.pattern.ask + import org.apache.pekko.pattern.ask import ClassicSimpleAggregate._ //All commands are valid in this test diff --git a/bounded-core/src/test/scala/io/cafienne/bounded/aggregate/ClassicSimpleAggregate.scala b/bounded-core/src/test/scala/io/cafienne/bounded/aggregate/ClassicSimpleAggregate.scala index 7057903..bb6084a 100644 --- a/bounded-core/src/test/scala/io/cafienne/bounded/aggregate/ClassicSimpleAggregate.scala +++ b/bounded-core/src/test/scala/io/cafienne/bounded/aggregate/ClassicSimpleAggregate.scala @@ -1,10 +1,10 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ package io.cafienne.bounded.aggregate -import akka.actor.{ActorSystem, Props} +import org.apache.pekko.actor.{ActorSystem, Props} import scala.collection.immutable.Seq import ClassicSimpleAggregate._ diff --git a/bounded-core/src/test/scala/io/cafienne/bounded/aggregate/TypedAnotherAggregate.scala b/bounded-core/src/test/scala/io/cafienne/bounded/aggregate/TypedAnotherAggregate.scala index f3b993c..c5d0f44 100644 --- a/bounded-core/src/test/scala/io/cafienne/bounded/aggregate/TypedAnotherAggregate.scala +++ b/bounded-core/src/test/scala/io/cafienne/bounded/aggregate/TypedAnotherAggregate.scala @@ -1,18 +1,18 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ package io.cafienne.bounded.aggregate -import akka.actor.typed.Behavior -import akka.cluster.sharding.typed.scaladsl.EntityTypeKey -import akka.persistence.typed.PersistenceId -import akka.persistence.typed.scaladsl.{Effect, EventSourcedBehavior, ReplyEffect} +import org.apache.pekko.actor.typed.Behavior +import org.apache.pekko.cluster.sharding.typed.scaladsl.EntityTypeKey +import org.apache.pekko.persistence.typed.PersistenceId +import org.apache.pekko.persistence.typed.scaladsl.{Effect, EventSourcedBehavior, ReplyEffect} import com.typesafe.scalalogging.Logger import io.cafienne.bounded.aggregate.typed.TypedAggregateRootManager import scala.concurrent.duration._ -import akka.actor.typed.scaladsl.{Behaviors, TimerScheduler} -import akka.persistence.RecoveryCompleted +import org.apache.pekko.actor.typed.scaladsl.{Behaviors, TimerScheduler} +import org.apache.pekko.persistence.RecoveryCompleted //This another aggregate is for testing the use of command gateways with multiple aggregates. object TypedAnotherAggregate { diff --git a/bounded-core/src/test/scala/io/cafienne/bounded/aggregate/TypedClusteredSpec.scala b/bounded-core/src/test/scala/io/cafienne/bounded/aggregate/TypedClusteredSpec.scala index 51b708b..ec44650 100644 --- a/bounded-core/src/test/scala/io/cafienne/bounded/aggregate/TypedClusteredSpec.scala +++ b/bounded-core/src/test/scala/io/cafienne/bounded/aggregate/TypedClusteredSpec.scala @@ -1,11 +1,11 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ package io.cafienne.bounded.aggregate import java.time.OffsetDateTime -import akka.actor.testkit.typed.scaladsl.ScalaTestWithActorTestKit +import org.apache.pekko.actor.testkit.typed.scaladsl.ScalaTestWithActorTestKit import org.scalatest.flatspec.AsyncFlatSpecLike import scala.concurrent.{ExecutionContextExecutor, Future} diff --git a/bounded-core/src/test/scala/io/cafienne/bounded/aggregate/TypedCommandGatewaySpec.scala b/bounded-core/src/test/scala/io/cafienne/bounded/aggregate/TypedCommandGatewaySpec.scala index d44732f..23c200f 100644 --- a/bounded-core/src/test/scala/io/cafienne/bounded/aggregate/TypedCommandGatewaySpec.scala +++ b/bounded-core/src/test/scala/io/cafienne/bounded/aggregate/TypedCommandGatewaySpec.scala @@ -1,26 +1,26 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ package io.cafienne.bounded.aggregate import java.time.OffsetDateTime -import akka.actor.testkit.typed.scaladsl.{ +import org.apache.pekko.actor.testkit.typed.scaladsl.{ LogCapturing, LoggingTestKit, ManualTime, ScalaTestWithActorTestKit, TestProbe } -import akka.actor.typed.{ActorSystem, Scheduler} -import akka.util.Timeout +import org.apache.pekko.actor.typed.{ActorSystem, Scheduler} +import org.apache.pekko.util.Timeout import io.cafienne.bounded.aggregate.typed.DefaultTypedCommandGateway import org.scalatest.BeforeAndAfterAll import org.scalatest.concurrent.ScalaFutures import org.scalatest.flatspec.AsyncFlatSpecLike import com.typesafe.config.ConfigFactory -import scala.concurrent.duration.* +import scala.concurrent.duration._ import scala.concurrent.{Await, ExecutionContextExecutor, Future} class TypedCommandGatewaySpec diff --git a/bounded-core/src/test/scala/io/cafienne/bounded/aggregate/TypedSimpleAggregate.scala b/bounded-core/src/test/scala/io/cafienne/bounded/aggregate/TypedSimpleAggregate.scala index d772b30..60624a0 100644 --- a/bounded-core/src/test/scala/io/cafienne/bounded/aggregate/TypedSimpleAggregate.scala +++ b/bounded-core/src/test/scala/io/cafienne/bounded/aggregate/TypedSimpleAggregate.scala @@ -1,20 +1,20 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ package io.cafienne.bounded.aggregate -import akka.actor.typed.{ActorRef, Behavior} -import akka.cluster.sharding.typed.scaladsl.EntityTypeKey -import akka.persistence.typed.PersistenceId -import akka.persistence.typed.scaladsl.{Effect, EventSourcedBehavior, ReplyEffect} +import org.apache.pekko.actor.typed.{ActorRef, Behavior} +import org.apache.pekko.cluster.sharding.typed.scaladsl.EntityTypeKey +import org.apache.pekko.persistence.typed.PersistenceId +import org.apache.pekko.persistence.typed.scaladsl.{Effect, EventSourcedBehavior, ReplyEffect} import com.typesafe.scalalogging.Logger import io.cafienne.bounded.aggregate.typed.TypedAggregateRootManager import java.time.OffsetDateTime import java.util.UUID import scala.concurrent.duration._ -import akka.actor.typed.scaladsl.{Behaviors, TimerScheduler} -import akka.persistence.RecoveryCompleted +import org.apache.pekko.actor.typed.scaladsl.{Behaviors, TimerScheduler} +import org.apache.pekko.persistence.RecoveryCompleted import io.cafienne.bounded.SampleProtocol.{CommandMetaData, MetaData, UserContext} object TypedSimpleAggregate { diff --git a/bounded-core/src/test/scala/io/cafienne/bounded/eventmaterializers/AbstractReplayableEventMaterializerWithRuntimeEventFilterSpec.scala b/bounded-core/src/test/scala/io/cafienne/bounded/eventmaterializers/AbstractReplayableEventMaterializerWithRuntimeEventFilterSpec.scala index 97f67ac..24fa863 100644 --- a/bounded-core/src/test/scala/io/cafienne/bounded/eventmaterializers/AbstractReplayableEventMaterializerWithRuntimeEventFilterSpec.scala +++ b/bounded-core/src/test/scala/io/cafienne/bounded/eventmaterializers/AbstractReplayableEventMaterializerWithRuntimeEventFilterSpec.scala @@ -1,172 +1,186 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ -package io.cafienne.bounded.eventmaterializers - -import java.time.OffsetDateTime -import akka.Done -import akka.actor.{ActorSystem, PoisonPill, Props} -import akka.event.{Logging, LoggingAdapter} -import akka.persistence.query.Sequence -import akka.testkit.{TestKit, TestProbe} -import akka.util.Timeout -import com.typesafe.scalalogging.Logger -import io.cafienne.bounded.aggregate.DomainEvent -import org.scalatest.BeforeAndAfterAll -import org.scalatest.concurrent.ScalaFutures -import org.scalatest.matchers.should.Matchers -import org.scalatest.time.{Millis, Seconds, Span} -import org.scalatest.wordspec.AnyWordSpec -import org.slf4j.LoggerFactory - -import scala.concurrent.Future -import scala.concurrent.duration._ - -case class TestMetaData( - timestamp: OffsetDateTime, - userContext: Option[String] -) - -case class TestedEvent(metaData: TestMetaData, text: String) extends DomainEvent { - def id: String = "entityId" -} - -class UserEventFilter( - userToFilter: String -) extends MaterializerEventFilter { - - override def filter(evt: DomainEvent): Boolean = { - evt match { - case TestedEvent(metaData, text) - if (metaData.userContext.isDefined && metaData.userContext.get.equalsIgnoreCase("user-b")) => - true - case _ => false - } - } - -} - -class AbstractReplayableEventMaterializerWithEventFilterSpec - extends AnyWordSpec - with Matchers - with ScalaFutures - with BeforeAndAfterAll { - - //Setup required supporting classes - implicit val timeout: Timeout = Timeout(10.seconds) - implicit val system: ActorSystem = ActorSystem("MaterializerTestSystem", SpecConfig.testConfig) - implicit val logger: LoggingAdapter = Logging(system, getClass) - implicit val defaultPatience: PatienceConfig = - PatienceConfig(timeout = Span(4, Seconds), interval = Span(100, Millis)) - - val eventStreamListener = TestProbe() - - val currentMeta = - TestMetaData(OffsetDateTime.parse("2018-01-01T17:43:00+01:00"), None) - - val testSet = Seq( - TestedEvent(currentMeta, "current-current"), - TestedEvent(currentMeta.copy(userContext = Some("user-a")), "current+1"), - TestedEvent(currentMeta.copy(userContext = Some("user-b")), "current+2"), - TestedEvent(currentMeta.copy(userContext = Some("user-a")), "current+3"), - TestedEvent(currentMeta.copy(userContext = Some("user-b")), "current+4"), - TestedEvent(currentMeta.copy(userContext = Some("user-b")), "current+5") - ) - - "The Event Materializer" must { - - "materialize all given events" in { - val materializer = new TestMaterializer() - - val toBeRun = new EventMaterializers(List(materializer)) - whenReady(toBeRun.startUp(false)) { replayResult => - logger.debug("replayResult: {}", replayResult) - assert(replayResult.head.offset == Some(Sequence(6L))) - } - logger.debug("DUMP all given events {}", materializer.storedEvents) - assert(materializer.storedEvents.size == 6) - } - - "materialize all events for user-b" in { - val materializer = new TestMaterializer(new UserEventFilter("user-b")) - - val toBeRun = new EventMaterializers(List(materializer)) - whenReady(toBeRun.startUp(false)) { replayResult => - logger.debug("replayResult: {}", replayResult) - assert(replayResult.head.offset == Some(Sequence(6L))) - } - logger.debug("DUMP current runtime and all versions {}", materializer.storedEvents) - assert(materializer.storedEvents.size == 3) - } - - } - - private def populateEventStore(evt: Seq[DomainEvent]): Unit = { - val storeEventsActor = system.actorOf(Props(classOf[CreateEventsInStoreActor], evt.head.id), "create-events-actor") - - val testProbe = TestProbe() - testProbe watch storeEventsActor - - system.eventStream.subscribe(eventStreamListener.ref, classOf[EventProcessed]) - - evt foreach { event => - testProbe.send(storeEventsActor, event) - testProbe.expectMsgAllConformingOf(classOf[DomainEvent]) - } - - storeEventsActor ! PoisonPill - val terminated = testProbe.expectTerminated(storeEventsActor) - assert(terminated.existenceConfirmed) - - } - - override def beforeAll(): Unit = { - populateEventStore(testSet) - } - - override def afterAll(): Unit = { - TestKit.shutdownActorSystem(system, 30.seconds, verifySystemShutdown = true) - } - - class TestMaterializer(eventFilter: MaterializerEventFilter = NoFilterEventFilter) - extends AbstractReplayableEventMaterializer( - system, - false, - eventFilter - ) { - - var storedEvents = Seq[DomainEvent]() - - override val logger: Logger = Logger(LoggerFactory.getLogger(TestMaterializer.this.getClass)) - - /** - * Tagname used to identify eventstream to listen to - */ - override val tagName: String = "testar" - - /** - * Mapping name of this listener - */ - override val matMappingName: String = "testar" - - /** - * Handle new incoming event - * - * @param evt event - */ - override def handleEvent(evt: Any): Future[Done] = { - logger.debug("TestMaterializer got event {} ", evt) - evt match { - case x: DomainEvent => storedEvents = storedEvents :+ x - case other => logger.warn("unkown event will not be stored {}", other) - } - Future.successful(Done) - } - - override def handleReplayEvent(evt: Any): Future[Done] = handleEvent(evt) - - override def toString: String = s"TestMaterializer $tagName contains ${storedEvents.mkString(",")}" - } - -} +///* +// * Copyright (C) 2016-2023 Batav B.V. +// */ +// +//package io.cafienne.bounded.eventmaterializers +// +//import java.time.OffsetDateTime +//import org.apache.pekko.Done +//import org.apache.pekko.actor.{ActorSystem, PoisonPill, Props} +//import org.apache.pekko.event.{Logging, LoggingAdapter} +//import org.apache.pekko.persistence.query.Sequence +//import org.apache.pekko.testkit.{TestKit, TestProbe} +//import org.apache.pekko.util.Timeout +//import com.typesafe.scalalogging.Logger +//import io.cafienne.bounded.aggregate.DomainEvent +//import org.apache.pekko.actor.testkit.typed.scaladsl.{LogCapturing, ScalaTestWithActorTestKit} +//import org.apache.pekko.persistence.Persistence +//import org.apache.pekko.persistence.testkit.PersistenceTestKitPlugin +//import org.apache.pekko.persistence.testkit.scaladsl.PersistenceTestKit +//import org.apache.pekko.projection.ProjectionId +//import org.apache.pekko.projection.testkit.scaladsl.{ProjectionTestKit, TestProjection, TestSourceProvider} +//import org.apache.pekko.stream.scaladsl.Source +//import org.scalatest.BeforeAndAfterAll +//import org.scalatest.concurrent.ScalaFutures +//import org.scalatest.matchers.should.Matchers +//import org.scalatest.time.{Millis, Seconds, Span} +//import org.scalatest.wordspec.{AnyWordSpec, AnyWordSpecLike} +//import org.slf4j.LoggerFactory +//import org.apache.pekko.projection.scaladsl.Handler +// +//import scala.concurrent.Future +//import scala.concurrent.duration.* +// +//case class TestMetaData( +// timestamp: OffsetDateTime, +// userContext: Option[String] +//) +// +//case class TestedEvent(metaData: TestMetaData, text: String) extends DomainEvent { +// def id: String = "entityId" +//} +// +//class UserEventFilter( +// userToFilter: String +//) extends MaterializerEventFilter { +// +// override def filter(evt: DomainEvent): Boolean = { +// evt match { +// case TestedEvent(metaData, text) +// if (metaData.userContext.isDefined && metaData.userContext.get.equalsIgnoreCase("user-b")) => +// true +// case _ => false +// } +// } +// +//} +// +//class AbstractReplayableEventMaterializerWithEventFilterSpec +// extends ScalaTestWithActorTestKit +// with AnyWordSpecLike +// with Matchers +// with ScalaFutures +// with LogCapturing { +// +// //Setup required supporting classes +// +//// implicit val system: ActorSystem = +//// ActorSystem("MaterializerTestSystem", PersistenceTestKitPlugin.config.withFallback(SpecConfig.testConfig)) +// // implicit val timeout: Timeout = Timeout(10.seconds) +// // implicit val logger: LoggingAdapter = Logging(system, getClass) +//// implicit val defaultPatience: PatienceConfig = +//// PatienceConfig(timeout = Span(4, Seconds), interval = Span(100, Millis)) +// +// //val testKit = PersistenceTestKit(system) +// private val projectionTestKit = ProjectionTestKit(system) +// val projectionId: ProjectionId = ProjectionId("name", "key") +//// val eventStreamListener = TestProbe() +// +// val currentMeta = +// TestMetaData(OffsetDateTime.parse("2018-01-01T17:43:00+01:00"), None) +// +// val testSet = Seq( +// TestedEvent(currentMeta, "current-current"), +// TestedEvent(currentMeta.copy(userContext = Some("user-a")), "current+1"), +// TestedEvent(currentMeta.copy(userContext = Some("user-b")), "current+2"), +// TestedEvent(currentMeta.copy(userContext = Some("user-a")), "current+3"), +// TestedEvent(currentMeta.copy(userContext = Some("user-b")), "current+4"), +// TestedEvent(currentMeta.copy(userContext = Some("user-b")), "current+5") +// ) +// +// def handler(strBuffer: StringBuffer, predicate: Int => Boolean): Handler[Int] = new Handler[Int] { +// override def process(env: Int): Future[Done] = { +// if (predicate(env)) concat(env) +// Future.successful(Done) +// } +// +// def concat(i: Int) = { +// if (strBuffer.toString.isEmpty) strBuffer.append(i) +// else strBuffer.append("-").append(i) +// } +// } +// +// "The Event Materializer" must { +// +// "run an function handler" in { +// val strBuffer = new StringBuffer() +// val sp = TestSourceProvider(Source(1 to 6), (i: Int) => i) +// val prj = TestProjection(projectionId, sp, () => handler(strBuffer, _ <= 6)) +// +// // stop as soon we observe that all expected elements passed through +// projectionTestKit.run(prj) { +// strBuffer.toString shouldBe "1-2-3-4-5-6" +// } +// } +// +//// "materialize all given events" in { +//// val materializer = new TestMaterializer() +//// +//// val toBeRun = new EventMaterializers(List(materializer)) +//// whenReady(toBeRun.startUp(false)) { replayResult => +//// system.log.debug("replayResult: {}", replayResult) +//// assert(replayResult.head.offset == Some(Sequence(6L))) +//// } +//// system.log.debug("DUMP all given events {}", materializer.storedEvents) +//// assert(materializer.storedEvents.size == 6) +//// } +// +//// "materialize all events for user-b" in { +//// val materializer = new TestMaterializer(new UserEventFilter("user-b")) +//// +//// val toBeRun = new EventMaterializers(List(materializer)) +//// whenReady(toBeRun.startUp(false)) { replayResult => +//// system.log.debug("replayResult: {}", replayResult) +//// assert(replayResult.head.offset == Some(Sequence(6L))) +//// } +//// system.log.debug("DUMP current runtime and all versions {}", materializer.storedEvents) +//// assert(materializer.storedEvents.size == 3) +//// } +// +// } +// +//// class TestMaterializer(eventFilter: MaterializerEventFilter = NoFilterEventFilter) +//// extends AbstractReplayableEventMaterializer( +//// system, +//// false, +//// eventFilter +//// ) { +//// +//// var storedEvents = Seq[DomainEvent]() +//// +//// override val logger: Logger = Logger(LoggerFactory.getLogger(TestMaterializer.this.getClass)) +//// +//// /** +//// * Tagname used to identify eventstream to listen to +//// */ +//// override val tagName: String = "testar" +//// +//// /** +//// * Mapping name of this listener +//// */ +//// override val matMappingName: String = "testar" +//// +//// /** +//// * Handle new incoming event +//// * +//// * @param evt event +//// */ +//// override def handleEvent(evt: Any): Future[Done] = { +//// logger.debug("TestMaterializer got event {} ", evt) +//// evt match { +//// case x: DomainEvent => storedEvents = storedEvents :+ x +//// case other => logger.warn("unkown event will not be stored {}", other) +//// } +//// Future.successful(Done) +//// } +//// +//// override def handleReplayEvent(evt: Any): Future[Done] = handleEvent(evt) +//// +//// override def toString: String = s"TestMaterializer $tagName contains ${storedEvents.mkString(",")}" +//// } +// +//} diff --git a/bounded-core/src/test/scala/io/cafienne/bounded/eventmaterializers/CreateEventsInStoreActor.scala b/bounded-core/src/test/scala/io/cafienne/bounded/eventmaterializers/CreateEventsInStoreActor.scala index b56aeb6..5e8b552 100644 --- a/bounded-core/src/test/scala/io/cafienne/bounded/eventmaterializers/CreateEventsInStoreActor.scala +++ b/bounded-core/src/test/scala/io/cafienne/bounded/eventmaterializers/CreateEventsInStoreActor.scala @@ -1,33 +1,37 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ -package io.cafienne.bounded.eventmaterializers - -import akka.actor.{ActorContext, ActorRef} -import akka.persistence.PersistentActor -import akka.persistence.journal.Tagged -import io.cafienne.bounded.aggregate.DomainEvent - -class CreateEventsInStoreActor(aggregateId: String) extends PersistentActor { - override def persistenceId: String = aggregateId - - override def receiveRecover: Receive = { - case other => - context.system.log.debug("received unknown event to recover:" + other) - } - - override def receiveCommand: Receive = { - case evt: Tagged => storeAndReply(sender(), evt) - case evt: DomainEvent => storeAndReply(sender(), evt) - case other => context.system.log.error(s"cannot handle storage of event $other") - } - - private def storeAndReply(replyTo: ActorRef, evt: Any)(implicit context: ActorContext): Unit = { - persist(evt) { e => - context.system.log.debug(s"persisted $e") - replyTo ! e - } - } - -} +///* +// * Copyright (C) 2016-2023 Batav B.V. +// */ +// +//package io.cafienne.bounded.eventmaterializers +// +//import org.apache.pekko.actor.{ActorContext, ActorRef} +//import org.apache.pekko.persistence.PersistentActor +//import org.apache.pekko.persistence.journal.Tagged +//import io.cafienne.bounded.aggregate.DomainEvent +// +//class CreateEventsInStoreActor(aggregateId: String) extends PersistentActor { +// override def persistenceId: String = aggregateId +// +// override def receiveRecover: Receive = { +// case other => +// context.system.log.debug("received unknown event to recover:" + other) +// } +// +// override def receiveCommand: Receive = { +// case evt: Tagged => storeAndReply(sender(), evt) +// case evt: DomainEvent => storeAndReply(sender(), evt) +// case other => context.system.log.error(s"cannot handle storage of event $other") +// } +// +// private def storeAndReply(replyTo: ActorRef, evt: Any)(implicit context: ActorContext): Unit = { +// persist(evt) { e => +// context.system.log.debug(s"persisted $e") +// replyTo ! e +// } +// } +// +//} diff --git a/bounded-core/src/test/scala/io/cafienne/bounded/eventmaterializers/SpecConfig.scala b/bounded-core/src/test/scala/io/cafienne/bounded/eventmaterializers/SpecConfig.scala index ec4f47e..9b5dd0c 100644 --- a/bounded-core/src/test/scala/io/cafienne/bounded/eventmaterializers/SpecConfig.scala +++ b/bounded-core/src/test/scala/io/cafienne/bounded/eventmaterializers/SpecConfig.scala @@ -1,77 +1,73 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ -package io.cafienne.bounded.eventmaterializers - -import com.typesafe.config.ConfigFactory - -object SpecConfig { - - /* - PLEASE NOTE: - Currently the https://github.com/dnvriend/akka-persistence-inmemory is NOT working for Aggregate Root tests - because it is not possible to use a separate instance writing the events that should be in the event store - before you actually create the aggregate root (should replay those stored events) to check execution of a new - command. - A new configuration that uses the akka bundled inmem storage is added to create a working situation. - */ - val testConfig = ConfigFactory.parseString( - """ - | akka { - | loglevel = "DEBUG" - | stdout-loglevel = "DEBUG" - | loggers = ["akka.testkit.TestEventListener"] - | actor { - | default-dispatcher { - | executor = "fork-join-executor" - | fork-join-executor { - | parallelism-min = 8 - | parallelism-factor = 2.0 - | parallelism-max = 8 - | } - | } - | serialize-creators = off - | serialize-messages = off - | serializers { - | //serializer = "io.cafienne.bounded.cargosample.persistence.CargoPersistersSerializer" - | } - | serialization-bindings { - | //"stamina.Persistable" = serializer - | // enable below to check if all events have been serialized without java.io.Serializable - | //"java.io.Serializable" = none - | } - | } - | persistence { - | publish-confirmations = on - | publish-plugin-commands = on - | journal { - | plugin = "inmemory-journal" - | } - | snapshot-store.plugin = "inmemory-snapshot-store" - | } - | test { - | single-expect-default = 10s - | timefactor = 1 - | } - | } - | inmemory-journal { - | event-adapters { - | testTagging = "io.cafienne.bounded.eventmaterializers.TestTaggingEventAdapter" - | } - | event-adapter-bindings { - | "io.cafienne.bounded.aggregate.DomainEvent" = testTagging - | } - | } - | inmemory-read-journal { - | refresh-interval = "10ms" - | max-buffer-size = "1000" - | } - | bounded.eventmaterializers.publish = true - | bounded.eventmaterializers.offsetstore { - | type = "inmemory" - | } - """.stripMargin - ) - -} +///* +// * Copyright (C) 2016-2023 Batav B.V. +// */ +// +//package io.cafienne.bounded.eventmaterializers +// +//import com.typesafe.config.ConfigFactory +// +//object SpecConfig { +// +// val testConfig = ConfigFactory.parseString( +// """ +// | pekko { +// | loglevel = "DEBUG" +// | stdout-loglevel = "DEBUG" +// | loggers = ["org.apache.pekko.testkit.TestEventListener"] +// | actor { +// | default-dispatcher { +// | executor = "fork-join-executor" +// | fork-join-executor { +// | parallelism-min = 8 +// | parallelism-factor = 2.0 +// | parallelism-max = 8 +// | } +// | } +// | serialize-creators = off +// | serialize-messages = off +// | serializers { +// | //serializer = "io.cafienne.bounded.cargosample.persistence.CargoPersistersSerializer" +// | } +// | serialization-bindings { +// | //"stamina.Persistable" = serializer +// | // enable below to check if all events have been serialized without java.io.Serializable +// | //"java.io.Serializable" = none +// | } +// | } +// | persistence { +// | publish-confirmations = on +// | publish-plugin-commands = on +// | journal { +// | plugin = "pekko.persistence.journal.inmem" +// | } +// | snapshot-store.plugin = "pekko.persistence.snapshot-store.local" +// | } +// | test { +// | single-expect-default = 10s +// | timefactor = 1 +// | } +// | } +// | inmemory-journal { +// | event-adapters { +// | testTagging = "io.cafienne.bounded.eventmaterializers.TestTaggingEventAdapter" +// | } +// | event-adapter-bindings { +// | "io.cafienne.bounded.aggregate.DomainEvent" = testTagging +// | } +// | } +// | inmemory-read-journal { +// | refresh-interval = "10ms" +// | max-buffer-size = "1000" +// | } +// | bounded.eventmaterializers.publish = true +// | bounded.eventmaterializers.offsetstore { +// | type = "inmemory" +// | } +// """.stripMargin +// ) +// +//} diff --git a/bounded-core/src/test/scala/io/cafienne/bounded/eventmaterializers/TestTaggingEventAdapter.scala b/bounded-core/src/test/scala/io/cafienne/bounded/eventmaterializers/TestTaggingEventAdapter.scala index 9979182..7f59307 100644 --- a/bounded-core/src/test/scala/io/cafienne/bounded/eventmaterializers/TestTaggingEventAdapter.scala +++ b/bounded-core/src/test/scala/io/cafienne/bounded/eventmaterializers/TestTaggingEventAdapter.scala @@ -1,18 +1,22 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ -package io.cafienne.bounded.eventmaterializers - -import akka.persistence.journal.{Tagged, WriteEventAdapter} -import io.cafienne.bounded.aggregate.DomainEvent - -class TestTaggingEventAdapter extends WriteEventAdapter { - override def manifest(event: Any): String = "" - - override def toJournal(event: Any): Any = event match { - case prEvent: DomainEvent => - Tagged(prEvent, Set("testar")) - case other => other - } -} +///* +// * Copyright (C) 2016-2023 Batav B.V. +// */ +// +//package io.cafienne.bounded.eventmaterializers +// +//import org.apache.pekko.persistence.journal.{Tagged, WriteEventAdapter} +//import io.cafienne.bounded.aggregate.DomainEvent +// +//class TestTaggingEventAdapter extends WriteEventAdapter { +// override def manifest(event: Any): String = "" +// +// override def toJournal(event: Any): Any = event match { +// case prEvent: DomainEvent => +// Tagged(prEvent, Set("testar")) +// case other => other +// } +//} diff --git a/bounded-pekko-http/src/main/scala/io/cafienne/bounded/akka/http/AkkaHttpParameterOffsetConverters.scala b/bounded-pekko-http/src/main/scala/io/cafienne/bounded/akka/http/AkkaHttpParameterOffsetConverters.scala index 0df8c65..906c61c 100644 --- a/bounded-pekko-http/src/main/scala/io/cafienne/bounded/akka/http/AkkaHttpParameterOffsetConverters.scala +++ b/bounded-pekko-http/src/main/scala/io/cafienne/bounded/akka/http/AkkaHttpParameterOffsetConverters.scala @@ -1,13 +1,13 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ package io.cafienne.bounded.akka.http import java.util.UUID -import akka.http.scaladsl.unmarshalling._ -import akka.persistence.query.Offset +import org.apache.pekko.http.scaladsl.unmarshalling._ +import org.apache.pekko.persistence.query.Offset import io.cafienne.bounded.config.Configured /** diff --git a/bounded-pekko-http/src/test/resources/logback-test.xml b/bounded-pekko-http/src/test/resources/logback-test.xml index 8fe094b..b62ee4e 100644 --- a/bounded-pekko-http/src/test/resources/logback-test.xml +++ b/bounded-pekko-http/src/test/resources/logback-test.xml @@ -7,7 +7,7 @@ - + diff --git a/bounded-test/src/main/scala/io/cafienne/bounded/akka/http/test/ScalatestTypedActorHttpRoute.scala b/bounded-test/src/main/scala/io/cafienne/bounded/akka/http/test/ScalatestTypedActorHttpRoute.scala index aeba5f1..b6e4697 100644 --- a/bounded-test/src/main/scala/io/cafienne/bounded/akka/http/test/ScalatestTypedActorHttpRoute.scala +++ b/bounded-test/src/main/scala/io/cafienne/bounded/akka/http/test/ScalatestTypedActorHttpRoute.scala @@ -1,18 +1,18 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ package io.cafienne.bounded.akka.http.test -import akka.actor.testkit.typed.scaladsl.{ActorTestKit, ActorTestKitBase} -import akka.actor.{ActorSystem, Scheduler} -import akka.http.scaladsl.testkit.ScalatestRouteTest -import akka.util.Timeout +import org.apache.pekko.actor.testkit.typed.scaladsl.{ActorTestKit, ActorTestKitBase} +import org.apache.pekko.actor.{ActorSystem, Scheduler} +import org.apache.pekko.http.scaladsl.testkit.ScalatestRouteTest +import org.apache.pekko.util.Timeout import org.scalatest.Suite //Thanks babloo80 -> https://github.com/akka/akka-http/issues/2036 trait ScalatestTypedActorHttpRoute extends ScalatestRouteTest { this: Suite => - import akka.actor.typed.scaladsl.adapter._ + import org.apache.pekko.actor.typed.scaladsl.adapter._ var typedTestKit : ActorTestKit = _ //val init causes createActorSystem() to cause NPE when typedTestKit.system is called in createActorSystem(). diff --git a/bounded-test/src/main/scala/io/cafienne/bounded/test/ClearStorageAfterEach.scala b/bounded-test/src/main/scala/io/cafienne/bounded/test/ClearStorageAfterEach.scala index b4a482a..87e244a 100644 --- a/bounded-test/src/main/scala/io/cafienne/bounded/test/ClearStorageAfterEach.scala +++ b/bounded-test/src/main/scala/io/cafienne/bounded/test/ClearStorageAfterEach.scala @@ -1,22 +1,23 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ -package io.cafienne.bounded.test - -import akka.persistence.inmemory.extension.{InMemoryJournalStorage, InMemorySnapshotStorage, StorageExtension} -import akka.testkit.{TestKit, TestProbe} -import org.scalatest.{BeforeAndAfterEach, Suite} - -trait ClearStorageAfterEach extends BeforeAndAfterEach { - this: TestKit with Suite => - - override protected def beforeEach(): Unit = { - val tp = TestProbe() - tp.send(StorageExtension(system).journalStorage, InMemoryJournalStorage.ClearJournal) - tp.expectMsg(akka.actor.Status.Success("")) - tp.send(StorageExtension(system).snapshotStorage, InMemorySnapshotStorage.ClearSnapshots) - tp.expectMsg(akka.actor.Status.Success("")) - super.beforeEach() - } -} +//package io.cafienne.bounded.test +// +////import org.apache.pekko.persistence.inmem.extension.{InMemoryJournalStorage, InMemorySnapshotStorage, StorageExtension} +//import org.apache.pekko.persistence.journal.inmem.InmemJournal +//import org.apache.pekko.testkit.{TestKit, TestProbe} +//import org.scalatest.{BeforeAndAfterEach, Suite} +// +//trait ClearStorageAfterEach extends BeforeAndAfterEach { +// this: TestKit with Suite => +// +// override protected def beforeEach(): Unit = { +// val tp = TestProbe() +// tp.send(StorageExtension(system).journalStorage, InMemoryJournalStorage.ClearJournal) +// tp.expectMsg(akka.actor.Status.Success("")) +// tp.send(StorageExtension(system).snapshotStorage, InMemorySnapshotStorage.ClearSnapshots) +// tp.expectMsg(akka.actor.Status.Success("")) +// super.beforeEach() +// } +//} diff --git a/bounded-test/src/main/scala/io/cafienne/bounded/test/CreateEventsInStoreActor.scala b/bounded-test/src/main/scala/io/cafienne/bounded/test/CreateEventsInStoreActor.scala index c41d833..b771ed0 100644 --- a/bounded-test/src/main/scala/io/cafienne/bounded/test/CreateEventsInStoreActor.scala +++ b/bounded-test/src/main/scala/io/cafienne/bounded/test/CreateEventsInStoreActor.scala @@ -1,11 +1,11 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ package io.cafienne.bounded.test -import akka.persistence.{PersistentActor, RecoveryCompleted} -import akka.persistence.journal.Tagged +import org.apache.pekko.persistence.{PersistentActor, RecoveryCompleted} +import org.apache.pekko.persistence.journal.Tagged import io.cafienne.bounded.aggregate.DomainEvent class CreateEventsInStoreActor(aggregateId: String, tags: Set[String] = Set.empty) extends PersistentActor { diff --git a/bounded-test/src/main/scala/io/cafienne/bounded/test/LoggingTestProbe.scala b/bounded-test/src/main/scala/io/cafienne/bounded/test/LoggingTestProbe.scala index 0712964..51b7a12 100644 --- a/bounded-test/src/main/scala/io/cafienne/bounded/test/LoggingTestProbe.scala +++ b/bounded-test/src/main/scala/io/cafienne/bounded/test/LoggingTestProbe.scala @@ -1,11 +1,11 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ package io.cafienne.bounded.test -import akka.actor.{ActorRef, ActorSystem} -import akka.testkit.{TestActor, TestProbe} +import org.apache.pekko.actor.{ActorRef, ActorSystem} +import org.apache.pekko.testkit.{TestActor, TestProbe} object LoggingTestProbe { def apply()(implicit system: ActorSystem): TestProbe = { diff --git a/bounded-test/src/main/scala/io/cafienne/bounded/test/StopSystemAfterAll.scala b/bounded-test/src/main/scala/io/cafienne/bounded/test/StopSystemAfterAll.scala index 53a877d..d83a669 100644 --- a/bounded-test/src/main/scala/io/cafienne/bounded/test/StopSystemAfterAll.scala +++ b/bounded-test/src/main/scala/io/cafienne/bounded/test/StopSystemAfterAll.scala @@ -1,10 +1,10 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ package io.cafienne.bounded.test -import akka.testkit.TestKit +import org.apache.pekko.testkit.TestKit import org.scalatest.{BeforeAndAfterAll, Suite} import scala.concurrent.duration._ diff --git a/bounded-test/src/main/scala/io/cafienne/bounded/test/TestableAggregateRoot.scala b/bounded-test/src/main/scala/io/cafienne/bounded/test/TestableAggregateRoot.scala index f970fcb..57c7f37 100644 --- a/bounded-test/src/main/scala/io/cafienne/bounded/test/TestableAggregateRoot.scala +++ b/bounded-test/src/main/scala/io/cafienne/bounded/test/TestableAggregateRoot.scala @@ -1,20 +1,23 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ package io.cafienne.bounded.test +import com.typesafe.config.ConfigFactory + import java.util.concurrent.atomic.AtomicInteger import io.cafienne.bounded.test.TestableAggregateRoot._ import scala.reflect.ClassTag -import akka.actor._ -import akka.pattern.ask -import akka.persistence.testkit.scaladsl.PersistenceTestKit -import akka.testkit.TestProbe -import akka.util.Timeout +import org.apache.pekko.actor._ +import org.apache.pekko.pattern.ask +import org.apache.pekko.persistence.testkit.scaladsl.PersistenceTestKit +import org.apache.pekko.testkit.TestProbe +import org.apache.pekko.util.Timeout import io.cafienne.bounded.aggregate.AggregateRootActor.GetState import io.cafienne.bounded.aggregate._ +import org.apache.pekko.persistence.testkit.PersistenceTestKitPlugin import scala.concurrent.Future import scala.concurrent.duration.Duration @@ -61,7 +64,7 @@ object TestableAggregateRoot { creator: AggregateRootCreator, id: String, evt: DomainEvent* - )(implicit system: ActorSystem, timeout: Timeout, ctag: reflect.ClassTag[A]): TestableAggregateRoot[A, B] = { + )(implicit timeout: Timeout, ctag: reflect.ClassTag[A]): TestableAggregateRoot[A, B] = { new TestableAggregateRoot[A, B](creator, id, evt) } @@ -75,7 +78,7 @@ object TestableAggregateRoot { def given[A <: AggregateRootActor[B], B <: AggregateState[B]: ClassTag]( creator: AggregateRootCreator, id: String - )(implicit system: ActorSystem, timeout: Timeout, ctag: reflect.ClassTag[A]): TestableAggregateRoot[A, B] = { + )(implicit timeout: Timeout, ctag: reflect.ClassTag[A]): TestableAggregateRoot[A, B] = { new TestableAggregateRoot[A, B](creator, id, Seq.empty[DomainEvent]) } @@ -91,10 +94,17 @@ class TestableAggregateRoot[A <: AggregateRootActor[B], B <: AggregateState[B]: id: String, evt: Seq[DomainEvent] )( - implicit system: ActorSystem, - timeout: Timeout, + implicit timeout: Timeout, ctag: reflect.ClassTag[A] ) { + implicit val system: ActorSystem = + ActorSystem( + "TestSystem", + PersistenceTestKitPlugin.config + .withFallback(ConfigFactory.parseString("akka.actor.allow-java-serialization=true")) + .withFallback(ConfigFactory.defaultApplication()) + .resolve() + ) implicit val duration: Duration = timeout.duration private var handledEvents: List[DomainEvent] = List.empty diff --git a/bounded-test/src/main/scala/io/cafienne/bounded/test/TestableProjection.scala b/bounded-test/src/main/scala/io/cafienne/bounded/test/TestableProjection.scala index b6c2c3c..716da1d 100644 --- a/bounded-test/src/main/scala/io/cafienne/bounded/test/TestableProjection.scala +++ b/bounded-test/src/main/scala/io/cafienne/bounded/test/TestableProjection.scala @@ -1,122 +1,127 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ -package io.cafienne.bounded.test - -import akka.Done - -import java.util.UUID -import akka.actor._ -import akka.persistence.inmemory.extension.{InMemoryJournalStorage, InMemorySnapshotStorage, StorageExtension} -import akka.testkit.TestProbe -import akka.util.Timeout -import io.cafienne.bounded.aggregate._ -import io.cafienne.bounded.eventmaterializers.{ - AbstractEventMaterializer, - EventMaterializers, - EventProcessed, - OffsetStoreProvider -} - -import scala.concurrent.duration._ -import scala.concurrent.{ExecutionContext, Future} -import akka.persistence.testkit.scaladsl.PersistenceTestKit - -object TestableProjection { - - def given( - evt: Seq[DomainEvent], - tags: Set[String] = Set.empty - )(implicit system: ActorSystem, timeout: Timeout = 2.seconds): TestableProjection = { - //Cleanup the store before this test is ran. - val tp = TestProbe() - tp.send(StorageExtension(system).journalStorage, InMemoryJournalStorage.ClearJournal) - tp.expectMsg(akka.actor.Status.Success("")) - tp.send(StorageExtension(system).snapshotStorage, InMemorySnapshotStorage.ClearSnapshots) - tp.expectMsg(akka.actor.Status.Success("")) - - val testedProjection = new TestableProjection(system, timeout, tags) - OffsetStoreProvider.getInMemoryStore().clear() - testedProjection.storeEvents(evt) - testedProjection - } - -} - -class TestableProjection private (system: ActorSystem, timeout: Timeout, tags: Set[String]) { - - private var eventMaterializers: Option[EventMaterializers] = _ - private implicit val executionContext: ExecutionContext = system.dispatcher - private implicit val actorSystem: ActorSystem = system - private var materializerId: Option[UUID] = None - - val eventStreamListener = TestProbe() -// val persistenceTestKit = PersistenceTestKit(system) - - if (!system.settings.config.hasPath("bounded.eventmaterializers.publish") || !system.settings.config.getBoolean( - "bounded.eventmaterializers.publish" - )) { - system.log.error("Config property bounded.eventmaterializers.publish must be enabled") - } - - def startProjection(projector: AbstractEventMaterializer): Future[EventMaterializers.ReplayResult] = { - materializerId = Some(projector.materializerId) - eventMaterializers = Some(new EventMaterializers(List(projector))) - eventMaterializers.get.startUp(true).map(list => list.head) - } - - def addEvents(evt: Seq[DomainEvent]): Future[Done] = { - eventMaterializers.fold(throw new IllegalStateException("You start the projection before you add events"))(_ => { - storeEvents(evt) - }) - } - - def addEvent(evt: DomainEvent): Future[Done] = { - eventMaterializers.fold(throw new IllegalStateException("You start the projection before you add events"))(_ => { - storeEvents(Seq(evt)) - }) - } - - // Blocking way to store events. - private def storeEvents(evt: Seq[DomainEvent]): Future[Done] = { - -// evt.groupBy(evt => evt.id).foreach(grp => persistenceTestKit.persistForRecovery(grp._1, grp._2)) - - val storeEventsActor = - system.actorOf(Props(classOf[CreateEventsInStoreActor], evt.head.id, tags), "create-events-actor") - - val testProbe = TestProbe() - testProbe watch storeEventsActor - - system.eventStream.subscribe(eventStreamListener.ref, classOf[EventProcessed]) - - evt foreach { event => - testProbe.send(storeEventsActor, event) - testProbe.expectMsgAllConformingOf(classOf[DomainEvent]) - } - - storeEventsActor ! PoisonPill - val terminated = testProbe.expectTerminated(storeEventsActor) - assert(terminated.existenceConfirmed) - - waitTillLastEventIsProcessed(evt) - } - - private def waitTillLastEventIsProcessed(evt: Seq[DomainEvent]): Future[Done] = { - if (materializerId.isDefined) { - var messageProcessedCounter = 0 - while (messageProcessedCounter < evt.length) { - eventStreamListener.fishForSpecificMessage(timeout.duration, "wait till last event is processed") { - case e: EventProcessed if materializerId.get == e.materializerId => - messageProcessedCounter += 1 - system.log.debug("catch " + e) - case other => system.log.debug("other message fished: {}", other) - } - } - Future.successful(Done) - } else { - Future.failed(new Exception("No materializer with id:" + materializerId + " found")) - } - } -} +///* +// * Copyright (C) 2016-2023 Batav B.V. +// */ +// +//package io.cafienne.bounded.test +// +//import org.apache.pekko.Done +// +//import java.util.UUID +//import org.apache.pekko.actor._ +////import org.apache.pekko.persistence.inmemory.extension.{InMemoryJournalStorage, InMemorySnapshotStorage, StorageExtension} +//import org.apache.pekko.testkit.TestProbe +//import org.apache.pekko.util.Timeout +//import io.cafienne.bounded.aggregate._ +//import io.cafienne.bounded.eventmaterializers.{ +// AbstractEventMaterializer, +// EventMaterializers, +// EventProcessed, +// OffsetStoreProvider +//} +// +//import scala.concurrent.duration._ +//import scala.concurrent.{ExecutionContext, Future} +//import org.apache.pekko.persistence.testkit.scaladsl.PersistenceTestKit +// +//object TestableProjection { +// +// def given( +// evt: Seq[DomainEvent], +// tags: Set[String] = Set.empty +// )(implicit system: ActorSystem, timeout: Timeout = 2.seconds): TestableProjection = { +// //Cleanup the store before this test is ran. +// val tp = TestProbe() +//// tp.send(StorageExtension(system).journalStorage, InMemoryJournalStorage.ClearJournal) +//// tp.expectMsg(akka.actor.Status.Success("")) +//// tp.send(StorageExtension(system).snapshotStorage, InMemorySnapshotStorage.ClearSnapshots) +//// tp.expectMsg(akka.actor.Status.Success("")) +// +// val testedProjection = new TestableProjection(system, timeout, tags) +// OffsetStoreProvider.getInMemoryStore().clear() +// testedProjection.storeEvents(evt) +// testedProjection +// } +// +//} +// +//class TestableProjection private (system: ActorSystem, timeout: Timeout, tags: Set[String]) { +// +// private var eventMaterializers: Option[EventMaterializers] = _ +// private implicit val executionContext: ExecutionContext = system.dispatcher +// private implicit val actorSystem: ActorSystem = system +// private var materializerId: Option[UUID] = None +// +// val eventStreamListener = TestProbe() +//// val persistenceTestKit = PersistenceTestKit(system) +// +// if (!system.settings.config.hasPath("bounded.eventmaterializers.publish") || !system.settings.config.getBoolean( +// "bounded.eventmaterializers.publish" +// )) { +// system.log.error("Config property bounded.eventmaterializers.publish must be enabled") +// } +// +// def startProjection(projector: AbstractEventMaterializer): Future[EventMaterializers.ReplayResult] = { +// materializerId = Some(projector.materializerId) +// eventMaterializers = Some(new EventMaterializers(List(projector))) +// eventMaterializers.get.startUp(true).map(list => list.head) +// } +// +// def addEvents(evt: Seq[DomainEvent]): Future[Done] = { +// eventMaterializers.fold(throw new IllegalStateException("You start the projection before you add events"))(_ => { +// storeEvents(evt) +// }) +// } +// +// def addEvent(evt: DomainEvent): Future[Done] = { +// eventMaterializers.fold(throw new IllegalStateException("You start the projection before you add events"))(_ => { +// storeEvents(Seq(evt)) +// }) +// } +// +// // Blocking way to store events. +// private def storeEvents(evt: Seq[DomainEvent]): Future[Done] = { +// +//// evt.groupBy(evt => evt.id).foreach(grp => persistenceTestKit.persistForRecovery(grp._1, grp._2)) +// +// val storeEventsActor = +// system.actorOf(Props(classOf[CreateEventsInStoreActor], evt.head.id, tags), "create-events-actor") +// +// val testProbe = TestProbe() +// testProbe watch storeEventsActor +// +// system.eventStream.subscribe(eventStreamListener.ref, classOf[EventProcessed]) +// +// evt foreach { event => +// testProbe.send(storeEventsActor, event) +// testProbe.expectMsgAllConformingOf(classOf[DomainEvent]) +// } +// +// storeEventsActor ! PoisonPill +// val terminated = testProbe.expectTerminated(storeEventsActor) +// //assert(terminated.existenceConfirmed) +// assert(terminated == Terminated) +// +// waitTillLastEventIsProcessed(evt) +// } +// +// private def waitTillLastEventIsProcessed(evt: Seq[DomainEvent]): Future[Done] = { +// if (materializerId.isDefined) { +// var messageProcessedCounter = 0 +// while (messageProcessedCounter < evt.length) { +// eventStreamListener.fishForSpecificMessage(timeout.duration, "wait till last event is processed") { +// case e: EventProcessed if materializerId.get == e.materializerId => +// messageProcessedCounter += 1 +// system.log.debug("catch " + e) +// case other => system.log.debug("other message fished: {}", other) +// } +// } +// Future.successful(Done) +// } else { +// Future.failed(new Exception("No materializer with id:" + materializerId + " found")) +// } +// } +//} diff --git a/bounded-test/src/main/scala/io/cafienne/bounded/test/TestableSaga.scala b/bounded-test/src/main/scala/io/cafienne/bounded/test/TestableSaga.scala index a68d813..2de7561 100644 --- a/bounded-test/src/main/scala/io/cafienne/bounded/test/TestableSaga.scala +++ b/bounded-test/src/main/scala/io/cafienne/bounded/test/TestableSaga.scala @@ -1,5 +1,5 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ package io.cafienne.bounded.test diff --git a/bounded-test/src/main/scala/io/cafienne/bounded/test/typed/TestTypedCommandGateway.scala b/bounded-test/src/main/scala/io/cafienne/bounded/test/typed/TestTypedCommandGateway.scala index 2981eda..c6bf7c1 100644 --- a/bounded-test/src/main/scala/io/cafienne/bounded/test/typed/TestTypedCommandGateway.scala +++ b/bounded-test/src/main/scala/io/cafienne/bounded/test/typed/TestTypedCommandGateway.scala @@ -1,12 +1,12 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ package io.cafienne.bounded.test.typed -import akka.actor.testkit.typed.scaladsl.ActorTestKit -import akka.actor.typed.{ActorRef, Scheduler} -import akka.util.Timeout +import org.apache.pekko.actor.testkit.typed.scaladsl.ActorTestKit +import org.apache.pekko.actor.typed.{ActorRef, Scheduler} +import org.apache.pekko.util.Timeout import io.cafienne.bounded.aggregate.typed.{TypedAggregateRootManager, TypedCommandGateway} import io.cafienne.bounded.aggregate.{DomainCommand, ValidateableCommand} @@ -21,7 +21,7 @@ class TestTypedCommandGateway[T <: DomainCommand]( ec: ExecutionContext, scheduler: Scheduler ) extends TypedCommandGateway[T] { - import akka.actor.typed.scaladsl.AskPattern._ + import org.apache.pekko.actor.typed.scaladsl.AskPattern._ val aggregates = mutable.Map[String, ActorRef[T]]() diff --git a/bounded-test/src/main/scala/io/cafienne/bounded/test/typed/TestableAggregateRoot.scala b/bounded-test/src/main/scala/io/cafienne/bounded/test/typed/TestableAggregateRoot.scala index 6d25882..753be9c 100644 --- a/bounded-test/src/main/scala/io/cafienne/bounded/test/typed/TestableAggregateRoot.scala +++ b/bounded-test/src/main/scala/io/cafienne/bounded/test/typed/TestableAggregateRoot.scala @@ -1,25 +1,28 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ package io.cafienne.bounded.test.typed +import com.typesafe.config.ConfigFactory + import java.util.concurrent.atomic.AtomicInteger -import akka.actor.{ActorSystem, typed} -import akka.actor.testkit.typed.scaladsl.ActorTestKit -import akka.actor.typed.eventstream.EventStream.Subscribe -import akka.actor.typed.{Behavior, ChildFailed, SupervisorStrategy} -import akka.actor.typed.scaladsl.Behaviors -import io.cafienne.bounded.test.typed.TestableAggregateRoot.* +import org.apache.pekko.actor.{ActorSystem, typed} +import org.apache.pekko.actor.testkit.typed.scaladsl.ActorTestKit +import org.apache.pekko.actor.typed.eventstream.EventStream.Subscribe +import org.apache.pekko.actor.typed.{Behavior, ChildFailed, SupervisorStrategy} +import org.apache.pekko.actor.typed.scaladsl.Behaviors +import io.cafienne.bounded.test.typed.TestableAggregateRoot._ import scala.reflect.ClassTag -import akka.util.Timeout +import org.apache.pekko.util.Timeout import io.cafienne.bounded.aggregate.{DomainCommand, DomainEvent, HandlingFailure} -import io.cafienne.bounded.aggregate.typed.* +import io.cafienne.bounded.aggregate.typed._ import scala.concurrent.duration.Duration -import akka.actor.typed.scaladsl.adapter.* -import akka.persistence.testkit.scaladsl.PersistenceTestKit +import org.apache.pekko.actor.typed.scaladsl.adapter._ +import org.apache.pekko.persistence.testkit.PersistenceTestKitPlugin +import org.apache.pekko.persistence.testkit.scaladsl.PersistenceTestKit import scala.collection.immutable @@ -67,7 +70,7 @@ object TestableAggregateRoot { creator: TypedAggregateRootManager[A], id: String, evt: B* - )(implicit system: ActorSystem, timeout: Timeout, ctag: reflect.ClassTag[A]): TestableAggregateRoot[A, B, C] = { + )(implicit timeout: Timeout, ctag: reflect.ClassTag[A]): TestableAggregateRoot[A, B, C] = { new TestableAggregateRoot[A, B, C](creator, id, immutable.Seq.concat(evt)) } @@ -81,7 +84,7 @@ object TestableAggregateRoot { def given[A <: DomainCommand, B <: DomainEvent, C: ClassTag]( creator: TypedAggregateRootManager[A], id: String - )(implicit system: ActorSystem, timeout: Timeout, ctag: reflect.ClassTag[A]): TestableAggregateRoot[A, B, C] = { + )(implicit timeout: Timeout, ctag: reflect.ClassTag[A]): TestableAggregateRoot[A, B, C] = { new TestableAggregateRoot[A, B, C](creator, id, immutable.Seq.empty[B]) } @@ -112,10 +115,17 @@ class TestableAggregateRoot[A <: DomainCommand, B <: DomainEvent, C: ClassTag] p evt: immutable.Seq[B], persistenceIdSeparator: String = "|" )( - implicit system: ActorSystem, - timeout: Timeout, + implicit timeout: Timeout, ctag: reflect.ClassTag[A] ) { + implicit val system: ActorSystem = + ActorSystem( + "TestSystem", + PersistenceTestKitPlugin.config + .withFallback(ConfigFactory.parseString("akka.actor.allow-java-serialization=on")) + .withFallback(ConfigFactory.defaultApplication()) + .resolve() + ) implicit val typedActorSystem: typed.ActorSystem[Nothing] = system.toTyped val testKit: ActorTestKit = ActorTestKit(system.toTyped) @@ -169,7 +179,9 @@ class TestableAggregateRoot[A <: DomainCommand, B <: DomainEvent, C: ClassTag] p if (command.aggregateRootId != id) throw MisdirectedCommand(id, command) wrappedActor.tell(command) lastCommand = Some(command) - Thread.sleep(1000) + val msgs = eventProbe.receiveMessages(1) + system.log.debug(s"found $msgs") + Thread.sleep(500) this } diff --git a/bounded-test/src/test/resources/logback-test.xml b/bounded-test/src/test/resources/logback-test.xml new file mode 100644 index 0000000..5c0638f --- /dev/null +++ b/bounded-test/src/test/resources/logback-test.xml @@ -0,0 +1,32 @@ + + + + + + INFO + + + [%date{ISO8601}] [%level] [%logger] [%marker] [%thread] - %msg MDC: {%mdc}%n + + + + + + + + + + + + + + + \ No newline at end of file diff --git a/bounded-test/src/test/resources/reference.conf b/bounded-test/src/test/resources/reference.conf new file mode 100644 index 0000000..8ed6055 --- /dev/null +++ b/bounded-test/src/test/resources/reference.conf @@ -0,0 +1,20 @@ +pekko { + loglevel = "DEBUG" + stdout-loglevel = "DEBUG" + loggers = ["org.apache.pekko.testkit.TestEventListener"] + actor { + allow-java-serialization = on + } + persistence { + publish-confirmations = on + publish-plugin-commands = on + journal { + plugin = "pekko.persistence.journal.inmem" + } + snapshot-store.plugin = "pekko.persistence.snapshot-store.local" + } + test { + single-expect-default = 10s + timefactor = 1 + } +} diff --git a/bounded-test/src/test/scala/io/cafienne/bounded/test/DomainProtocol.scala b/bounded-test/src/test/scala/io/cafienne/bounded/test/DomainProtocol.scala index 25dfc58..1d46b38 100644 --- a/bounded-test/src/test/scala/io/cafienne/bounded/test/DomainProtocol.scala +++ b/bounded-test/src/test/scala/io/cafienne/bounded/test/DomainProtocol.scala @@ -1,5 +1,5 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ package io.cafienne.bounded.test diff --git a/bounded-test/src/test/scala/io/cafienne/bounded/test/SpecConfig.scala b/bounded-test/src/test/scala/io/cafienne/bounded/test/SpecConfig.scala index 4065b6e..603ba6d 100644 --- a/bounded-test/src/test/scala/io/cafienne/bounded/test/SpecConfig.scala +++ b/bounded-test/src/test/scala/io/cafienne/bounded/test/SpecConfig.scala @@ -1,5 +1,5 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ package io.cafienne.bounded.test @@ -18,10 +18,10 @@ object SpecConfig { */ val testConfig = ConfigFactory.parseString( """ - | akka { + | pekko { | loglevel = "DEBUG" | stdout-loglevel = "DEBUG" - | loggers = ["akka.testkit.TestEventListener"] + | loggers = ["org.apache.pekko.testkit.TestEventListener"] | actor { | serialize-messages = off | serialize-creators = off @@ -41,19 +41,15 @@ object SpecConfig { | publish-confirmations = on | publish-plugin-commands = on | journal { - | plugin = "inmemory-journal" + | plugin = "pekko.persistence.journal.inmem" | } - | snapshot-store.plugin = "inmemory-snapshot-store" + | snapshot-store.plugin = "pekko.persistence.snapshot-store.local" | } | test { | single-expect-default = 10s | timefactor = 1 | } | } - | inmemory-read-journal { - | refresh-interval = "10ms" - | max-buffer-size = "1000" - | } | | bounded.eventmaterializers.publish = true | diff --git a/bounded-test/src/test/scala/io/cafienne/bounded/test/TestAggregateRoot.scala b/bounded-test/src/test/scala/io/cafienne/bounded/test/TestAggregateRoot.scala index 9687bed..1a96b8a 100644 --- a/bounded-test/src/test/scala/io/cafienne/bounded/test/TestAggregateRoot.scala +++ b/bounded-test/src/test/scala/io/cafienne/bounded/test/TestAggregateRoot.scala @@ -1,10 +1,10 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ package io.cafienne.bounded.test -import akka.actor.{ActorSystem, Props} +import org.apache.pekko.actor.{ActorSystem, Props} import io.cafienne.bounded.aggregate._ import io.cafienne.bounded.test.DomainProtocol.StateUpdated import io.cafienne.bounded.test.TestAggregateRoot.TestAggregateRootState diff --git a/bounded-test/src/test/scala/io/cafienne/bounded/test/TestableAggregateRootSpec.scala b/bounded-test/src/test/scala/io/cafienne/bounded/test/TestableAggregateRootSpec.scala index 0e5b36f..014c047 100644 --- a/bounded-test/src/test/scala/io/cafienne/bounded/test/TestableAggregateRootSpec.scala +++ b/bounded-test/src/test/scala/io/cafienne/bounded/test/TestableAggregateRootSpec.scala @@ -1,16 +1,16 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ package io.cafienne.bounded.test import java.time.OffsetDateTime -import akka.actor.ActorSystem -import akka.actor.testkit.typed.scaladsl.ScalaTestWithActorTestKit -import akka.event.{Logging, LoggingAdapter} -import akka.persistence.testkit.PersistenceTestKitPlugin -import akka.testkit.TestKit -import akka.util.Timeout +import org.apache.pekko.actor.ActorSystem +import org.apache.pekko.actor.testkit.typed.scaladsl.ScalaTestWithActorTestKit +import org.apache.pekko.event.{Logging, LoggingAdapter} +import org.apache.pekko.persistence.testkit.PersistenceTestKitPlugin +import org.apache.pekko.testkit.TestKit +import org.apache.pekko.util.Timeout import io.cafienne.bounded.test.TestableAggregateRoot._ import io.cafienne.bounded.test.DomainProtocol._ import io.cafienne.bounded.test.TestAggregateRoot.TestAggregateRootState diff --git a/bounded-test/src/test/scala/io/cafienne/bounded/test/TestableProjectionSpec.scala b/bounded-test/src/test/scala/io/cafienne/bounded/test/TestableProjectionSpec.scala index 4a47515..0c88a9f 100644 --- a/bounded-test/src/test/scala/io/cafienne/bounded/test/TestableProjectionSpec.scala +++ b/bounded-test/src/test/scala/io/cafienne/bounded/test/TestableProjectionSpec.scala @@ -1,125 +1,129 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ -package io.cafienne.bounded.test - -import akka.actor.ActorSystem -import akka.event.{Logging, LoggingAdapter} -import akka.util.Timeout -import org.scalatest.BeforeAndAfterAll -import org.scalatest.concurrent.ScalaFutures -import org.scalatest.matchers.should.Matchers -import org.scalatest.wordspec.AsyncWordSpecLike -import akka.Done -import com.typesafe.scalalogging.Logger -import io.cafienne.bounded.eventmaterializers.{AbstractReplayableEventMaterializer, OffsetStoreProvider} -import org.slf4j.LoggerFactory - -import scala.concurrent.{ExecutionContext, Future} -import scala.concurrent.duration._ -import DomainProtocol._ -import akka.persistence.query.Sequence -import io.cafienne.bounded.aggregate.DomainEvent - -import java.time.OffsetDateTime - -class TestableProjectionSpec extends AsyncWordSpecLike with Matchers with ScalaFutures with BeforeAndAfterAll { - - implicit val timeout: Timeout = Timeout(30.seconds) - implicit val system: ActorSystem = - ActorSystem("TestableProjectionSpecSystem", SpecConfig.testConfig) - - implicit val logger: LoggingAdapter = Logging(system, getClass) - - val metaDate: TestMetaData = TestMetaData(OffsetDateTime.now(), None, None) - - "The testable projection" must { - val testProjectionMaterializer = new TestProjectionMaterializer(system) with OffsetStoreProvider - val testAggregateRootId1 = "arTest1" - val evt1 = InitialStateCreated(metaDate, testAggregateRootId1, "initialState") - val evt2 = StateUpdated(metaDate, testAggregateRootId1, "updatedState") - val fixture = TestableProjection.given(Seq(evt1, evt2), Set("ar-test", "aggregate")) - - "Store basic events at the start" in { - whenReady(fixture.startProjection(testProjectionMaterializer)) { replayResult => - logger.info("replayResult: {}", replayResult) - assert(replayResult.offset == Some(Sequence(2L))) - } - } - - "Store more events" in { - val evt3 = StateUpdated(metaDate, testAggregateRootId1, "morestate") - whenReady(fixture.addEvent(evt3)) { result => testProjectionMaterializer.events.size should be(3) } - } - - "Store more events slowly" in { - val evt4 = SlowStateUpdated(metaDate, testAggregateRootId1, "morestate", 1000L) - whenReady(fixture.addEvent(evt4)) { result => testProjectionMaterializer.events.size should be(4) } - } - - "Store even more events slowly" in { - val evt5 = SlowStateUpdated(metaDate, testAggregateRootId1, "morestate2", 1000L) - val evt6 = SlowStateUpdated(metaDate, testAggregateRootId1, "morestate3", 1000L) - whenReady(fixture.addEvents(Seq(evt5, evt6))) { result => testProjectionMaterializer.events.size should be(6) } - } - - "Store even more and more events slowly" in { - val evt7 = SlowStateUpdated(metaDate, testAggregateRootId1, "morestate4", 2000L) - val evt8 = SlowStateUpdated(metaDate, testAggregateRootId1, "morestate5", 2000L) - val evt9 = SlowStateUpdated(metaDate, testAggregateRootId1, "morestate6", 2000L) - whenReady(fixture.addEvents(Seq(evt7, evt8, evt9))) { result => - testProjectionMaterializer.events.size should be(9) - } - } - } - -} - -class TestProjectionMaterializer(actorSystem: ActorSystem) extends AbstractReplayableEventMaterializer(actorSystem) { - - /** - * Tagname used to identify eventstream to listen to - */ - override val tagName: String = "ar-test" - - /** - * Mapping name of this listener - */ - override val matMappingName: String = "test-view" - - override lazy val logger: Logger = Logger(LoggerFactory.getLogger(TestProjectionMaterializer.this.getClass)) - - implicit val ec: ExecutionContext = system.dispatcher - - var events: Seq[DomainEvent] = Seq.empty[DomainEvent] - - override def handleReplayEvent(evt: Any): Future[Done] = handleEvent(evt) - - override def handleEvent(evt: Any): Future[Done] = { - try { - evt match { - case event: InitialStateCreated => - events = events :+ event - Future.successful(Done) - case event: StateUpdated => - events = events :+ event - Future.successful(Done) - case event: SlowStateUpdated => - Thread.sleep(event.waited) - events = events :+ event - Future.successful(Done) - case _ => Future.successful(Done) - } - } catch { - case ex: Throwable => - logger.error( - "Unable to process event: " + evt.getClass.getSimpleName + Option(ex.getCause) - .map(ex => ex.getMessage) - .getOrElse("") + s" ${ex.getMessage} " + " exception: " + logException(ex), - ex - ) - Future.failed(ex) - } - } -} +///* +// * Copyright (C) 2016-2023 Batav B.V. +// */ +// +//package io.cafienne.bounded.test +// +//import org.apache.pekko.actor.ActorSystem +//import org.apache.pekko.event.{Logging, LoggingAdapter} +//import org.apache.pekko.util.Timeout +//import org.scalatest.BeforeAndAfterAll +//import org.scalatest.concurrent.ScalaFutures +//import org.scalatest.matchers.should.Matchers +//import org.scalatest.wordspec.AsyncWordSpecLike +//import org.apache.pekko.Done +//import com.typesafe.scalalogging.Logger +//import io.cafienne.bounded.eventmaterializers.{AbstractReplayableEventMaterializer, OffsetStoreProvider} +//import org.slf4j.LoggerFactory +// +//import scala.concurrent.{ExecutionContext, Future} +//import scala.concurrent.duration._ +//import DomainProtocol._ +//import org.apache.pekko.persistence.query.Sequence +//import io.cafienne.bounded.aggregate.DomainEvent +// +//import java.time.OffsetDateTime +// +//class TestableProjectionSpec extends AsyncWordSpecLike with Matchers with ScalaFutures with BeforeAndAfterAll { +// +// implicit val timeout: Timeout = Timeout(30.seconds) +// implicit val system: ActorSystem = +// ActorSystem("TestableProjectionSpecSystem", SpecConfig.testConfig) +// +// implicit val logger: LoggingAdapter = Logging(system, getClass) +// +// val metaDate: TestMetaData = TestMetaData(OffsetDateTime.now(), None, None) +// +// "The testable projection" must { +// val testProjectionMaterializer = new TestProjectionMaterializer(system) with OffsetStoreProvider +// val testAggregateRootId1 = "arTest1" +// val evt1 = InitialStateCreated(metaDate, testAggregateRootId1, "initialState") +// val evt2 = StateUpdated(metaDate, testAggregateRootId1, "updatedState") +// val fixture = TestableProjection.given(Seq(evt1, evt2), Set("ar-test", "aggregate")) +// +// "Store basic events at the start" in { +// whenReady(fixture.startProjection(testProjectionMaterializer)) { replayResult => +// logger.info("replayResult: {}", replayResult) +// assert(replayResult.offset == Some(Sequence(2L))) +// } +// } +// +// "Store more events" in { +// val evt3 = StateUpdated(metaDate, testAggregateRootId1, "morestate") +// whenReady(fixture.addEvent(evt3)) { result => testProjectionMaterializer.events.size should be(3) } +// } +// +// "Store more events slowly" in { +// val evt4 = SlowStateUpdated(metaDate, testAggregateRootId1, "morestate", 1000L) +// whenReady(fixture.addEvent(evt4)) { result => testProjectionMaterializer.events.size should be(4) } +// } +// +// "Store even more events slowly" in { +// val evt5 = SlowStateUpdated(metaDate, testAggregateRootId1, "morestate2", 1000L) +// val evt6 = SlowStateUpdated(metaDate, testAggregateRootId1, "morestate3", 1000L) +// whenReady(fixture.addEvents(Seq(evt5, evt6))) { result => testProjectionMaterializer.events.size should be(6) } +// } +// +// "Store even more and more events slowly" in { +// val evt7 = SlowStateUpdated(metaDate, testAggregateRootId1, "morestate4", 2000L) +// val evt8 = SlowStateUpdated(metaDate, testAggregateRootId1, "morestate5", 2000L) +// val evt9 = SlowStateUpdated(metaDate, testAggregateRootId1, "morestate6", 2000L) +// whenReady(fixture.addEvents(Seq(evt7, evt8, evt9))) { result => +// testProjectionMaterializer.events.size should be(9) +// } +// } +// } +// +//} +// +//class TestProjectionMaterializer(actorSystem: ActorSystem) extends AbstractReplayableEventMaterializer(actorSystem) { +// +// /** +// * Tagname used to identify eventstream to listen to +// */ +// override val tagName: String = "ar-test" +// +// /** +// * Mapping name of this listener +// */ +// override val matMappingName: String = "test-view" +// +// override lazy val logger: Logger = Logger(LoggerFactory.getLogger(TestProjectionMaterializer.this.getClass)) +// +// implicit val ec: ExecutionContext = system.dispatcher +// +// var events: Seq[DomainEvent] = Seq.empty[DomainEvent] +// +// override def handleReplayEvent(evt: Any): Future[Done] = handleEvent(evt) +// +// override def handleEvent(evt: Any): Future[Done] = { +// try { +// evt match { +// case event: InitialStateCreated => +// events = events :+ event +// Future.successful(Done) +// case event: StateUpdated => +// events = events :+ event +// Future.successful(Done) +// case event: SlowStateUpdated => +// Thread.sleep(event.waited) +// events = events :+ event +// Future.successful(Done) +// case _ => Future.successful(Done) +// } +// } catch { +// case ex: Throwable => +// logger.error( +// "Unable to process event: " + evt.getClass.getSimpleName + Option(ex.getCause) +// .map(ex => ex.getMessage) +// .getOrElse("") + s" ${ex.getMessage} " + " exception: " + logException(ex), +// ex +// ) +// Future.failed(ex) +// } +// } +//} diff --git a/bounded-test/src/test/scala/io/cafienne/bounded/test/typed/TestableAggregateRootSpec.scala b/bounded-test/src/test/scala/io/cafienne/bounded/test/typed/TestableAggregateRootSpec.scala index 54f0c69..5d718ad 100644 --- a/bounded-test/src/test/scala/io/cafienne/bounded/test/typed/TestableAggregateRootSpec.scala +++ b/bounded-test/src/test/scala/io/cafienne/bounded/test/typed/TestableAggregateRootSpec.scala @@ -1,22 +1,22 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ package io.cafienne.bounded.test.typed -import akka.actor.{ActorSystem, typed} +import org.apache.pekko.actor.{ActorSystem, typed} import java.time.OffsetDateTime -import akka.actor.testkit.typed.scaladsl.{ActorTestKit, ScalaTestWithActorTestKit, TestProbe} -import akka.util.Timeout +import org.apache.pekko.actor.testkit.typed.scaladsl.{ActorTestKit, ScalaTestWithActorTestKit, TestProbe} +import org.apache.pekko.util.Timeout import io.cafienne.bounded.aggregate.typed.TypedAggregateRootManager -import io.cafienne.bounded.test.DomainProtocol.* +import io.cafienne.bounded.test.DomainProtocol._ import org.scalatest.BeforeAndAfterAll import org.scalatest.concurrent.ScalaFutures import org.scalatest.flatspec.AsyncFlatSpecLike -import akka.actor.typed.scaladsl.adapter.* -import akka.persistence.testkit.PersistenceTestKitPlugin -import akka.persistence.testkit.scaladsl.PersistenceTestKit +import org.apache.pekko.actor.typed.scaladsl.adapter._ +import org.apache.pekko.persistence.testkit.PersistenceTestKitPlugin +import org.apache.pekko.persistence.testkit.scaladsl.PersistenceTestKit import com.typesafe.config.ConfigFactory import io.cafienne.bounded.test.typed.TestableAggregateRoot.{ AggregateThrownException, @@ -24,41 +24,24 @@ import io.cafienne.bounded.test.typed.TestableAggregateRoot.{ NoCommandsIssued, UnexpectedCommandHandlingSuccess } +import org.scalatest.matchers.should.Matchers -import scala.concurrent.duration.* +import scala.concurrent.duration._ -class TestableAggregateRootSpec - extends ScalaTestWithActorTestKit( - PersistenceTestKitPlugin.config - .withFallback(ConfigFactory.parseString("akka.actor.allow-java-serialization=true")) - .withFallback(ConfigFactory.defaultApplication()) - ) - with AsyncFlatSpecLike - with ScalaFutures - with BeforeAndAfterAll { //with LogCapturing { +class TestableAggregateRootSpec extends AsyncFlatSpecLike with Matchers with ScalaFutures { //with LogCapturing { import TypedSimpleAggregate._ - implicit val classicActorSystem: ActorSystem = system.toClassic - implicit val typedTestKit: ActorTestKit = testKit - val persistenceTestKit: PersistenceTestKit = PersistenceTestKit(system) behavior of "Testable Typed Aggregate Root" - implicit val gatewayTimeout: Timeout = Timeout(10.seconds) - implicit val actorSytem: typed.ActorSystem[Nothing] = system - + implicit val gatewayTimeout: Timeout = Timeout(10.seconds) val creator: TypedAggregateRootManager[SimpleAggregateCommand] = new SimpleAggregateManager() -// override def beforeEach(): Unit = { -// persistenceTestKit.clearAll() -// } - "TestableAggregateRoot" should "be created with initial events" in { val commandMetaData = TestCommandMetaData(OffsetDateTime.now(), None) val eventMetaData = TestMetaData.fromCommand(commandMetaData) val aggregateId = "ar0" - val testProbe = TestProbe[Response]() val testAR = TestableAggregateRoot .given( creator, @@ -66,10 +49,11 @@ class TestableAggregateRootSpec Created(aggregateId, eventMetaData), ItemAdded(aggregateId, eventMetaData, "ar0 item 1") ) - .when(AddItem(aggregateId, commandMetaData, "ar0 item 2", testProbe.ref)) - testProbe.expectMessage(OK) - testAR.events.size should be(1) + testAR.when(AddItem(aggregateId, commandMetaData, "ar0 item 2")) //, testProbe.ref)) + + //testProbe.expectMessage(OK) + //testAR.events.size should be(1) //Initial events are not shown in the list, so you know what events were generated by the command testAR.events should be(List(ItemAdded(aggregateId, eventMetaData, "ar0 item 2"))) } @@ -78,7 +62,6 @@ class TestableAggregateRootSpec val commandMetaData = TestCommandMetaData(OffsetDateTime.now(), None) val eventMetaData = TestMetaData.fromCommand(commandMetaData) val testAggregateRootId1 = "3" - val testProbe = TestProbe[Response]() TestableAggregateRoot .given( @@ -86,7 +69,7 @@ class TestableAggregateRootSpec testAggregateRootId1, Created(testAggregateRootId1, eventMetaData) ) - .when(TriggerError(testAggregateRootId1, commandMetaData, testProbe.ref)) + .when(TriggerError(testAggregateRootId1, commandMetaData)) .failure match { case AggregateThrownException(ex: IllegalArgumentException) => ex.getMessage should be("This is an exception thrown during processing the command for AR 3") @@ -98,21 +81,18 @@ class TestableAggregateRootSpec val commandMetaData = TestCommandMetaData(OffsetDateTime.now(), None) val aggregateRootId = "4" val wrongId = "5" - val testProbe = TestProbe[Response]() an[MisdirectedCommand] should be thrownBy { TestableAggregateRoot .given(creator, aggregateRootId) - .when(Create(wrongId, commandMetaData, testProbe.ref)) + .when(Create(wrongId, commandMetaData)) } } it should "signal that handling was successful despite expected failure" in { val commandMetaData = TestCommandMetaData(OffsetDateTime.now(), None) val eventMetaData = TestMetaData.fromCommand(commandMetaData) - val aggregateRootId = "3" - val testProbe = TestProbe[Response]() an[UnexpectedCommandHandlingSuccess] should be thrownBy { TestableAggregateRoot @@ -121,7 +101,7 @@ class TestableAggregateRootSpec aggregateRootId, Created(aggregateRootId, eventMetaData) ) - .when(AddItem(aggregateRootId, commandMetaData, "new item", testProbe.ref)) + .when(AddItem(aggregateRootId, commandMetaData, "new item")) .failure } //TODO ... strange ? @@ -178,8 +158,4 @@ class TestableAggregateRootSpec // // } - protected override def afterAll(): Unit = { - ActorTestKit.shutdown(system, 30.seconds, throwIfShutdownFails = true) - } - } diff --git a/bounded-test/src/test/scala/io/cafienne/bounded/test/typed/TypedSimpleAggregate.scala b/bounded-test/src/test/scala/io/cafienne/bounded/test/typed/TypedSimpleAggregate.scala index f4862e3..a593db7 100644 --- a/bounded-test/src/test/scala/io/cafienne/bounded/test/typed/TypedSimpleAggregate.scala +++ b/bounded-test/src/test/scala/io/cafienne/bounded/test/typed/TypedSimpleAggregate.scala @@ -1,15 +1,15 @@ /* - * Copyright (C) 2016-2023 Batav B.V. + * Copyright (C) 2016-2024 Batav B.V. */ package io.cafienne.bounded.test.typed -import akka.actor.typed.scaladsl.{Behaviors, TimerScheduler} -import akka.actor.typed.{ActorRef, Behavior} -import akka.cluster.sharding.typed.scaladsl.EntityTypeKey -import akka.persistence.RecoveryCompleted -import akka.persistence.typed.PersistenceId -import akka.persistence.typed.scaladsl.{Effect, EventSourcedBehavior, ReplyEffect} +import org.apache.pekko.actor.typed.scaladsl.{Behaviors, TimerScheduler} +import org.apache.pekko.actor.typed.{ActorRef, Behavior} +import org.apache.pekko.cluster.sharding.typed.scaladsl.EntityTypeKey +import org.apache.pekko.persistence.RecoveryCompleted +import org.apache.pekko.persistence.typed.PersistenceId +import org.apache.pekko.persistence.typed.scaladsl.{Effect, EventSourcedBehavior, ReplyEffect} import com.typesafe.scalalogging.Logger import io.cafienne.bounded.aggregate.{DomainCommand, DomainEvent} import io.cafienne.bounded.aggregate.typed.TypedAggregateRootManager @@ -30,14 +30,12 @@ object TypedSimpleAggregate { sealed trait SimpleAggregateCommand extends DomainCommand - final case class Create(aggregateRootId: String, metaData: TestCommandMetaData, replyTo: ActorRef[Response]) - extends SimpleAggregateCommand + final case class Create(aggregateRootId: String, metaData: TestCommandMetaData) extends SimpleAggregateCommand final case class AddItem( aggregateRootId: String, metaData: TestCommandMetaData, - item: String, - replyTo: ActorRef[Response] + item: String ) extends SimpleAggregateCommand final case class Stop(aggregateRootId: String, metaData: TestCommandMetaData, replyTo: ActorRef[Response]) @@ -53,8 +51,7 @@ object TypedSimpleAggregate { replyTo: ActorRef[Response] ) extends SimpleAggregateCommand - final case class TriggerError(aggregateRootId: String, metaData: TestCommandMetaData, replyTo: ActorRef[Response]) - extends SimpleAggregateCommand + final case class TriggerError(aggregateRootId: String, metaData: TestCommandMetaData) extends SimpleAggregateCommand //NOTE that a GET on the aggregate is not according to the CQRS pattern and added here for testing. sealed trait SimpleDirectAggregateQuery extends SimpleAggregateCommand @@ -113,14 +110,14 @@ object TypedSimpleAggregate { private def createAggregate(cmd: Create): ReplyEffect[SimpleAggregateEvent, SimpleAggregateState] = { logger.debug(s"Create Aggregate (replayed: $replayed)" + cmd) - Effect.persist(Created(cmd.aggregateRootId, TestMetaData.fromCommand(cmd.metaData))).thenReply(cmd.replyTo)(_ ⇒ OK) + Effect.persist(Created(cmd.aggregateRootId, TestMetaData.fromCommand(cmd.metaData))).thenNoReply() } private def addItem(cmd: AddItem): ReplyEffect[SimpleAggregateEvent, SimpleAggregateState] = { logger.debug("AddItem " + cmd) Effect .persist(ItemAdded(cmd.aggregateRootId, TestMetaData.fromCommand(cmd.metaData), cmd.item)) - .thenReply(cmd.replyTo)(_ ⇒ OK) + .thenNoReply() } // event handler to keep internal aggregate state diff --git a/build.sbt b/build.sbt index ee017ba..83352df 100644 --- a/build.sbt +++ b/build.sbt @@ -1,6 +1,6 @@ lazy val basicSettings = { - val scala213 = "2.13.11" + val scala213 = "2.13.14" val supportedScalaVersions = List(scala213) Seq( @@ -18,11 +18,11 @@ lazy val basicSettings = { "-Xlint", // recommended additional warnings "-Ywarn-value-discard", // Warn when non-Unit expression results are unused "-Ywarn-dead-code", - "-Ywarn-unused", - "-Xsource:3" //, + "-Ywarn-unused" //, + //"-Xsource:3" //, //"-Ywarn-unused-import" ), - scalastyleConfig := baseDirectory.value / "project/scalastyle-config.xml", + //scalastyleConfig := baseDirectory.value / "project/scalastyle-config.xml", scalafmtConfig := (ThisBuild / baseDirectory).value / "project/.scalafmt.conf", scalafmtOnCompile := true, @@ -46,15 +46,24 @@ lazy val basicSettings = { url("https://github.com/olger")) ), // Add sonatype repository settings - publishTo := Some( - if (isSnapshot.value) - Opts.resolver.sonatypeSnapshots - else - Opts.resolver.sonatypeStaging - ), + publishTo := { + // For accounts created after Feb 2021: + // val nexus = "https://s01.oss.sonatype.org/" + val nexus = "https://oss.sonatype.org/" + if (isSnapshot.value) Some("snapshots" at nexus + "content/repositories/snapshots") + else Some("releases" at nexus + "service/local/staging/deploy/maven2") + }, +// publishTo := Some( +// if (isSnapshot.value) { +// Opts.resolver.sonatypeOssSnapshots +// } else { +// Opts.resolver.sonatypeStaging +// } +// ), publishMavenStyle := true, Test / publishArtifact := false, pomIncludeRepository := { _ => false }, + resolvers += Resolver.ApacheMavenSnapshotsRepo, ) } @@ -72,7 +81,8 @@ val boundedCore = (project in file("bounded-core")) .settings(basicSettings: _*) .settings( name := "bounded-core", - libraryDependencies ++= Dependencies.baseDeps ++ Dependencies.persistanceLmdbDBDeps ++ Dependencies.persistenceCassandraDeps ++ Dependencies.persistenceJdbcDeps ++ Dependencies.testDeps) + //libraryDependencies ++= Dependencies.baseDeps ++ Dependencies.persistanceLmdbDBDeps ++ Dependencies.persistenceCassandraDeps ++ Dependencies.persistenceJdbcDeps ++ Dependencies.testDeps) + libraryDependencies ++= Dependencies.baseDeps ++ Dependencies.persistenceLevelDBDeps ++ Dependencies.persistenceCassandraDeps ++ Dependencies.persistenceJdbcDeps ++ Dependencies.testDeps) val boundedAkkaHttp = (project in file("bounded-pekko-http")) .dependsOn(boundedCore) diff --git a/project/Dependencies.scala b/project/Dependencies.scala index d545d08..f948dd0 100644 --- a/project/Dependencies.scala +++ b/project/Dependencies.scala @@ -6,39 +6,40 @@ import sbt._ object Dependencies { - val akkaVersion = "2.6.21" - val akkaHttpVersion = "10.2.9" - val staminaVersion = "0.1.6" - val persistenceInMemVersion = "2.5.15.2" - val scalaTestVersion = "3.2.11" + val pekkoVersion = "1.1.2" + val pekkoHttpVersion = "1.1.0" + val pekkoProjectionSnapshotVersion = "1.1.0-M1" + val scalaTestVersion = "3.2.19" val baseDeps = { - def akkaModule(name: String, version: String = akkaVersion) = - "com.typesafe.akka" %% s"akka-$name" % version + def pekkoModule(name: String, version: String = pekkoVersion) = + "org.apache.pekko" %% s"pekko-$name" % version Seq( - akkaModule("slf4j"), - akkaModule("actor"), - akkaModule("stream"), - akkaModule("persistence"), - akkaModule("persistence-typed"), - akkaModule("persistence-query"), - akkaModule("cluster-sharding-typed"), - akkaModule("cluster"), - akkaModule("coordination"), - akkaModule("cluster-tools"), - akkaModule("stream-testkit") % Test, - akkaModule("testkit") % Test, - akkaModule("actor-testkit-typed") % Test, - akkaModule("persistence-testkit") % Test, + pekkoModule("slf4j"), + pekkoModule("actor"), + pekkoModule("stream"), + pekkoModule("persistence"), + pekkoModule("persistence-typed"), + pekkoModule("persistence-query"), + pekkoModule("projection-core", pekkoProjectionSnapshotVersion), + pekkoModule("projection-eventsourced", pekkoProjectionSnapshotVersion), + pekkoModule("cluster-sharding-typed"), + pekkoModule("cluster"), + pekkoModule("coordination"), + pekkoModule("cluster-tools"), + pekkoModule("stream-testkit") % Test, + pekkoModule("testkit") % Test, + pekkoModule("actor-testkit-typed") % Test, + pekkoModule("persistence-testkit") % Test, + pekkoModule("projection-testkit",pekkoProjectionSnapshotVersion) % Test, "io.spray" %% "spray-json" % "1.3.6", - "com.github.dnvriend" %% "akka-persistence-inmemory" % persistenceInMemVersion, "com.typesafe.scala-logging" %% "scala-logging" % "3.9.5" ) } val log = Seq( - "ch.qos.logback" % "logback-classic" % "1.4.7", - "net.logstash.logback" % "logstash-logback-encoder" % "7.3" + "ch.qos.logback" % "logback-classic" % "1.5.12", + "net.logstash.logback" % "logstash-logback-encoder" % "8.0" ) @@ -48,8 +49,8 @@ object Dependencies { val akkaHttpDeps = { - def akkaHttpModule(name: String, version: String = akkaHttpVersion) = - "com.typesafe.akka" %% s"akka-$name" % version + def akkaHttpModule(name: String, version: String = pekkoHttpVersion) = + "org.apache.pekko" %% s"pekko-$name" % version baseDeps ++ Seq( akkaHttpModule("http"), @@ -61,10 +62,10 @@ object Dependencies { val testDeps = { baseDeps ++ Seq( "org.scalatest" %% "scalatest" % scalaTestVersion, - "com.typesafe.akka" %% "akka-testkit" % akkaVersion, - "com.typesafe.akka" %% "akka-actor-testkit-typed" % akkaVersion, - "com.typesafe.akka" %% "akka-persistence-testkit" % akkaVersion, - "com.typesafe.akka" %% "akka-http-testkit" % akkaHttpVersion + "org.apache.pekko" %% "pekko-testkit" % pekkoVersion, + "org.apache.pekko" %% "pekko-actor-testkit-typed" % pekkoVersion, + "org.apache.pekko" %% "pekko-persistence-testkit" % pekkoVersion, + "org.apache.pekko" %% "pekko-http-testkit" % pekkoHttpVersion ) ++ test } @@ -77,20 +78,21 @@ object Dependencies { val persistanceLmdbDBDeps = { baseDeps ++ Seq( - "org.lmdbjava" % "lmdbjava" % "0.8.2" + "org.lmdbjava" % "lmdbjava" % "0.9.0" ) } val persistenceCassandraDeps = { baseDeps ++ Seq( - "com.typesafe.akka" %% "akka-persistence-cassandra" % "1.0.5" + "org.apache.pekko" %% "pekko-persistence-cassandra" % "1.1.0-M1" ) } val persistenceJdbcDeps = { baseDeps ++ Seq( - "com.lightbend.akka" %% "akka-persistence-jdbc" % "5.0.4", - "com.typesafe.slick" %% "slick" % "3.3.3" + "org.apache.pekko" %% "pekko-persistence-r2dbc" % "1.0.0", + "org.apache.pekko" %% "pekko-persistence-jdbc" % "1.1.0", + "com.typesafe.slick" %% "slick" % "3.5.2", ) } diff --git a/project/build.properties b/project/build.properties index 3d7427d..db1723b 100644 --- a/project/build.properties +++ b/project/build.properties @@ -1 +1 @@ -sbt.version=1.9.2 \ No newline at end of file +sbt.version=1.10.5 diff --git a/project/plugins.sbt b/project/plugins.sbt index 4c86e71..4985a8b 100644 --- a/project/plugins.sbt +++ b/project/plugins.sbt @@ -1,11 +1,11 @@ -addSbtPlugin("org.scalameta" % "sbt-scalafmt" % "2.3.4") -addSbtPlugin("org.scoverage" % "sbt-scoverage" % "1.6.1") -addSbtPlugin("org.scalastyle" % "scalastyle-sbt-plugin" % "1.0.0") -addSbtPlugin("org.xerial.sbt" % "sbt-sonatype" % "2.3") +addSbtPlugin("org.scalameta" % "sbt-scalafmt" % "2.4.6") +addSbtPlugin("org.scoverage" % "sbt-scoverage" % "2.2.2") +//addSbtPlugin("org.scalastyle" % "scalastyle-sbt-plugin" % "1.0.0") +addSbtPlugin("org.xerial.sbt" % "sbt-sonatype" % "3.9.11") addSbtPlugin("com.jsuereth" % "sbt-pgp" % "1.1.0") addSbtPlugin("de.heikoseeberger" % "sbt-header" % "5.6.0") addSbtPlugin("com.timushev.sbt" % "sbt-updates" % "0.5.0") -addSbtPlugin("com.github.gseitz" % "sbt-release" % "1.0.13") -addSbtPlugin("com.eed3si9n" % "sbt-buildinfo" % "0.9.0") +addSbtPlugin("com.github.sbt" % "sbt-release" % "1.1.0") +addSbtPlugin("com.eed3si9n" % "sbt-buildinfo" % "0.10.0") addSbtPlugin("com.typesafe.sbt" % "sbt-git" % "1.0.0") -addSbtPlugin("net.virtual-void" % "sbt-dependency-graph" % "0.10.0-RC1") +//addSbtPlugin("net.virtual-void" % "sbt-dependency-graph" % "0.10.0-RC1") diff --git a/version.sbt b/version.sbt index 40bc6a2..8eb9c89 100644 --- a/version.sbt +++ b/version.sbt @@ -1,2 +1,2 @@ -ThisBuild / version := "0.4.0" +ThisBuild / version := "0.4.0-SNAPSHOT" From 77be397624e631d592ca9be127509d4bfa72509b Mon Sep 17 00:00:00 2001 From: Olger Warnier Date: Tue, 26 Nov 2024 20:41:34 +0100 Subject: [PATCH 3/5] Update Sonar --- .circleci/config.yml | 5 +---- 1 file changed, 1 insertion(+), 4 deletions(-) diff --git a/.circleci/config.yml b/.circleci/config.yml index b519384..ed08359 100644 --- a/.circleci/config.yml +++ b/.circleci/config.yml @@ -41,13 +41,10 @@ jobs: - "~/.ivy2/cache" - "~/.sbt" - "~/.m2" - - run: - name: Store Coverage Report - command: bash <(curl -s https://codecov.io/bash) - sonarcloud/scan orbs: - sonarcloud: sonarsource/sonarcloud@1.1.1 + sonarcloud: sonarsource/sonarcloud@2.0.0 workflows: version: 2 From bf0a796e1e852bf4014f1065ce0f7adbae22a94f Mon Sep 17 00:00:00 2001 From: Olger Warnier Date: Tue, 26 Nov 2024 20:50:48 +0100 Subject: [PATCH 4/5] Update scala --- build.sbt | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/build.sbt b/build.sbt index 83352df..888b64c 100644 --- a/build.sbt +++ b/build.sbt @@ -1,6 +1,6 @@ lazy val basicSettings = { - val scala213 = "2.13.14" + val scala213 = "2.13.15" val supportedScalaVersions = List(scala213) Seq( From 677e47e7d874422d15a2f55ca8d56cde1c19104b Mon Sep 17 00:00:00 2001 From: Olger Warnier Date: Tue, 26 Nov 2024 20:53:51 +0100 Subject: [PATCH 5/5] Akka to pekko for sonar --- sonar-project.properties | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/sonar-project.properties b/sonar-project.properties index ca63e28..f1a7b22 100644 --- a/sonar-project.properties +++ b/sonar-project.properties @@ -9,7 +9,7 @@ sonar.projectName=Bounded Framework #sonar.projectVersion=1.0 # Path is relative to the sonar-project.properties file. Defaults to . -sonar.sources=bounded-core/src,bounded-test/src,bounded-akka-http/src +sonar.sources=bounded-core/src,bounded-test/src,bounded-pekko-http/src # Encoding of the source code. Default is default system encoding sonar.sourceEncoding=UTF-8