Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 1 addition & 4 deletions .circleci/config.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
@@ -1,11 +1,11 @@
/*
* Copyright (C) 2016-2023 Batav B.V. <https://www.cafienne.io/bounded>
* Copyright (C) 2016-2024 Batav B.V. <https://www.cafienne.io/bounded>
*/

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 {
Expand Down
Original file line number Diff line number Diff line change
@@ -1,13 +1,13 @@
/*
* Copyright (C) 2016-2023 Batav B.V. <https://www.cafienne.io/bounded>
* Copyright (C) 2016-2024 Batav B.V. <https://www.cafienne.io/bounded>
*/

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}
Expand All @@ -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)
Expand Down Expand Up @@ -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)
Expand All @@ -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
Expand Down
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
/*
* Copyright (C) 2016-2023 Batav B.V. <https://www.cafienne.io/bounded>
* Copyright (C) 2016-2024 Batav B.V. <https://www.cafienne.io/bounded>
*/

package io.cafienne.bounded.aggregate
Expand Down
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
/*
* Copyright (C) 2016-2023 Batav B.V. <https://www.cafienne.io/bounded>
* Copyright (C) 2016-2024 Batav B.V. <https://www.cafienne.io/bounded>
*/

package io.cafienne.bounded.aggregate
Expand Down
Original file line number Diff line number Diff line change
@@ -1,10 +1,10 @@
/*
* Copyright (C) 2016-2023 Batav B.V. <https://www.cafienne.io/bounded>
* Copyright (C) 2016-2024 Batav B.V. <https://www.cafienne.io/bounded>
*/

package io.cafienne.bounded.aggregate

import akka.actor.typed.ActorRef
import org.apache.pekko.actor.typed.ActorRef
import scala.collection.immutable.Seq

trait DomainCommand {
Expand Down
Original file line number Diff line number Diff line change
@@ -1,12 +1,12 @@
/*
* Copyright (C) 2016-2023 Batav B.V. <https://www.cafienne.io/bounded>
* Copyright (C) 2016-2024 Batav B.V. <https://www.cafienne.io/bounded>
*/

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}
Expand Down Expand Up @@ -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 {

Expand Down
Original file line number Diff line number Diff line change
@@ -1,11 +1,11 @@
/*
* Copyright (C) 2016-2023 Batav B.V. <https://www.cafienne.io/bounded>
* Copyright (C) 2016-2024 Batav B.V. <https://www.cafienne.io/bounded>
*/

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] {
Expand Down
Original file line number Diff line number Diff line change
@@ -1,10 +1,10 @@
/*
* Copyright (C) 2016-2023 Batav B.V. <https://www.cafienne.io/bounded>
* Copyright (C) 2016-2024 Batav B.V. <https://www.cafienne.io/bounded>
*/

package io.cafienne.bounded.akka

import akka.actor.ActorSystem
import org.apache.pekko.actor.ActorSystem

trait ActorSystemProvider {
implicit def system: ActorSystem
Expand Down
Original file line number Diff line number Diff line change
@@ -1,28 +1,27 @@
/*
* Copyright (C) 2016-2023 Batav B.V. <https://www.cafienne.io/bounded>
* Copyright (C) 2016-2024 Batav B.V. <https://www.cafienne.io/bounded>
*/

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 = {
Expand All @@ -32,38 +31,28 @@ 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
]
}
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
]
Expand Down
Original file line number Diff line number Diff line change
@@ -1,11 +1,11 @@
/*
* Copyright (C) 2016-2023 Batav B.V. <https://www.cafienne.io/bounded>
* Copyright (C) 2016-2024 Batav B.V. <https://www.cafienne.io/bounded>
*/

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 {

Expand Down
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
/*
* Copyright (C) 2016-2023 Batav B.V. <https://www.cafienne.io/bounded>
* Copyright (C) 2016-2024 Batav B.V. <https://www.cafienne.io/bounded>
*/

package io.cafienne.bounded.config
Expand Down
Original file line number Diff line number Diff line change
@@ -1,15 +1,15 @@
/*
* Copyright (C) 2016-2023 Batav B.V. <https://www.cafienne.io/bounded>
* Copyright (C) 2016-2024 Batav B.V. <https://www.cafienne.io/bounded>
*/

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
Expand All @@ -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
Expand Down
Original file line number Diff line number Diff line change
@@ -1,13 +1,13 @@
/*
* Copyright (C) 2016-2023 Batav B.V. <https://www.cafienne.io/bounded>
* Copyright (C) 2016-2024 Batav B.V. <https://www.cafienne.io/bounded>
*/

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
Expand Down
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
/*
* Copyright (C) 2016-2023 Batav B.V. <https://www.cafienne.io/bounded>
* Copyright (C) 2016-2024 Batav B.V. <https://www.cafienne.io/bounded>
*/

package io.cafienne.bounded.eventmaterializers
Expand Down
Original file line number Diff line number Diff line change
@@ -1,10 +1,10 @@
/*
* Copyright (C) 2016-2023 Batav B.V. <https://www.cafienne.io/bounded>
* Copyright (C) 2016-2024 Batav B.V. <https://www.cafienne.io/bounded>
*/

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

Expand Down
Original file line number Diff line number Diff line change
@@ -1,10 +1,10 @@
/*
* Copyright (C) 2016-2023 Batav B.V. <https://www.cafienne.io/bounded>
* Copyright (C) 2016-2024 Batav B.V. <https://www.cafienne.io/bounded>
*/

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)
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
/*
* Copyright (C) 2016-2023 Batav B.V. <https://www.cafienne.io/bounded>
* Copyright (C) 2016-2024 Batav B.V. <https://www.cafienne.io/bounded>
*/

package io.cafienne.bounded.eventmaterializers
Expand Down
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
/*
* Copyright (C) 2016-2023 Batav B.V. <https://www.cafienne.io/bounded>
* Copyright (C) 2016-2024 Batav B.V. <https://www.cafienne.io/bounded>
*/

package io.cafienne.bounded.eventmaterializers
Expand Down
Original file line number Diff line number Diff line change
@@ -1,11 +1,11 @@
/*
* Copyright (C) 2016-2023 Batav B.V. <https://www.cafienne.io/bounded>
* Copyright (C) 2016-2024 Batav B.V. <https://www.cafienne.io/bounded>
*/

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.{
Expand Down
Loading