This is an automated email from the ASF dual-hosted git repository.
pjfanning pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/pekko-projection.git
The following commit(s) were added to refs/heads/main by this push:
new 18f73ba1 chore: remove leftover Scala 2.12 idioms (#625)
18f73ba1 is described below
commit 18f73ba1581f1d4ac2d40ec5e8eceb4db1f89e7a
Author: PJ Fanning <[email protected]>
AuthorDate: Sun Sep 13 08:37:49 2026 +0100
chore: remove leftover Scala 2.12 idioms (#625)
Motivation:
The codebase still carried the Scala 2.12-era `immutable.Seq` qualifier
throughout main, test and doc sources. `scala.Seq` has been an alias for
`scala.collection.immutable.Seq` since Scala 2.13, so the qualifier and the
`scala.collection.immutable` imports that exist only to support it are
redundant. This mirrors apache/pekko#3539.
Modification:
- Replace `immutable.Seq` with `Seq` across 37 Scala files and drop the
`scala.collection.immutable` imports that are no longer referenced
(kept where `immutable.IndexedSeq`/`immutable.Set` are still used).
- Drop the explicit `import scala.collection.immutable.Seq` lines.
- Paradox prose (cassandra, jdbc, r2dbc, slick) updated from
`immutable.Seq` to `Seq` to match the doc snippets.
The other idioms addressed upstream (`WrappedArray`, `filterKeys`/
`mapValues`, `Either` projections, `toIterator`, `Stream`) do not occur in
this repository.
Result:
No remaining Scala 2.12 collection idioms in main, test or doc sources.
Source- and binary-compatible: same type, same erasure.
Tests:
- native `scalafmt --mode diff-ref=upstream/main` run on changed files
- `sbt Test/compile jdbc-int-test/Test/compile r2dbc-int-test/Test/compile
slick-int-test/Test/compile` passes on Scala 2.13 and 3.3
- `sbt core/mimaReportBinaryIssues jdbc/mimaReportBinaryIssues
slick/mimaReportBinaryIssues cassandra/mimaReportBinaryIssues
eventsourced/mimaReportBinaryIssues durable-state/mimaReportBinaryIssues
r2dbc/mimaReportBinaryIssues` passes (grpc has MiMa disabled)
- `git diff --check` clean
References:
None - mirrors apache/pekko#3539; Scala 2.12 support was dropped previously
---
.../cassandra/CassandraProjectionSpec.scala | 5 ++--
.../cassandra/scaladsl/CassandraProjection.scala | 3 +--
.../internal/OffsetSerializationSpec.scala | 6 ++---
.../tools/InternalProjectionStateMetricsSpec.scala | 5 ++--
.../internal/metrics/tools/TestHandlers.scala | 7 +++--
.../pekko/projection/internal/HandlerAdapter.scala | 7 +++--
.../internal/InternalProjectionState.scala | 15 +++++------
.../projection/internal/OffsetSerialization.scala | 4 +--
.../pekko/projection/internal/OffsetStrategy.scala | 5 ++--
.../pekko/projection/internal/Telemetry.scala | 3 +--
docs/src/main/paradox/cassandra.md | 2 +-
docs/src/main/paradox/jdbc.md | 2 +-
docs/src/main/paradox/r2dbc.md | 2 +-
docs/src/main/paradox/slick.md | 2 +-
.../scaladsl/DurableStateSourceProvider.scala | 3 +--
.../scaladsl/EventSourcedProvider.scala | 5 ++--
.../scaldsl/EventSourcedProviderSpec.scala | 1 -
.../scala/docs/jdbc/JdbcProjectionDocExample.scala | 6 ++---
.../docs/slick/SlickProjectionDocExample.scala | 6 ++---
.../grpc/internal/EventProducerServiceSpec.scala | 3 +--
.../projection/grpc/consumer/ConsumerFilter.scala | 29 ++++++++++-----------
.../grpc/consumer/scaladsl/GrpcReadJournal.scala | 13 +++++-----
.../grpc/internal/ConsumerFilterRegistry.scala | 7 +++--
.../grpc/internal/ConsumerFilterStore.scala | 13 +++++-----
.../grpc/internal/ProtoAnySerialization.scala | 3 +--
.../cassandra/CassandraProjectionDocExample.scala | 6 ++---
.../pekko/projection/jdbc/JdbcProjectionSpec.scala | 7 +++--
.../pekko/projection/jdbc/internal/Dialect.scala | 30 ++++++++++------------
.../jdbc/internal/JdbcHandlerAdapter.scala | 7 +++--
.../jdbc/internal/JdbcProjectionImpl.scala | 7 +++--
.../projection/jdbc/scaladsl/JdbcProjection.scala | 5 ++--
.../integration/KafkaToSlickIntegrationSpec.scala | 4 +--
.../projection/r2dbc/R2dbcProjectionSpec.scala | 7 +++--
.../r2dbc/R2dbcTimestampOffsetProjectionSpec.scala | 24 ++++++++---------
.../r2dbc/internal/R2dbcHandlerAdapter.scala | 7 +++--
.../r2dbc/internal/R2dbcOffsetStore.scala | 4 +--
.../r2dbc/internal/R2dbcProjectionImpl.scala | 15 +++++------
.../r2dbc/scaladsl/R2dbcProjection.scala | 10 +++-----
.../projection/R2dbcProjectionDocExample.scala | 6 ++---
.../pekko/projection/slick/SlickProjection.scala | 11 ++++----
.../projection/slick/SlickProjectionSpec.scala | 7 +++--
41 files changed, 135 insertions(+), 179 deletions(-)
diff --git
a/cassandra-test/src/test/scala/org/apache/pekko/projection/cassandra/CassandraProjectionSpec.scala
b/cassandra-test/src/test/scala/org/apache/pekko/projection/cassandra/CassandraProjectionSpec.scala
index 41c3ea15..f2989afd 100644
---
a/cassandra-test/src/test/scala/org/apache/pekko/projection/cassandra/CassandraProjectionSpec.scala
+++
b/cassandra-test/src/test/scala/org/apache/pekko/projection/cassandra/CassandraProjectionSpec.scala
@@ -18,7 +18,6 @@ import java.util.concurrent.atomic.AtomicInteger
import java.util.concurrent.atomic.AtomicReference
import scala.annotation.tailrec
-import scala.collection.immutable
import scala.concurrent.Await
import scala.concurrent.ExecutionContext
import scala.concurrent.Future
@@ -584,10 +583,10 @@ class CassandraProjectionSpec
val entityId = UUID.randomUUID().toString
val projectionId = genRandomProjectionId()
- def groupedHandler(): Handler[immutable.Seq[Envelope]] = new
Handler[immutable.Seq[Envelope]] {
+ def groupedHandler(): Handler[Seq[Envelope]] = new
Handler[Seq[Envelope]] {
private var state: Future[Option[ConcatStr]] =
repository.findById(entityId)
- override def process(group: immutable.Seq[Envelope]): Future[Done] = {
+ override def process(group: Seq[Envelope]): Future[Done] = {
val newState = state.flatMap { s =>
val concatStr = group.foldLeft(s) {
case (None, env) => Some(ConcatStr(env.id, env.message))
diff --git
a/cassandra/src/main/scala/org/apache/pekko/projection/cassandra/scaladsl/CassandraProjection.scala
b/cassandra/src/main/scala/org/apache/pekko/projection/cassandra/scaladsl/CassandraProjection.scala
index 5f74fb0c..c25b6c64 100644
---
a/cassandra/src/main/scala/org/apache/pekko/projection/cassandra/scaladsl/CassandraProjection.scala
+++
b/cassandra/src/main/scala/org/apache/pekko/projection/cassandra/scaladsl/CassandraProjection.scala
@@ -13,7 +13,6 @@
package org.apache.pekko.projection.cassandra.scaladsl
-import scala.collection.immutable
import scala.concurrent.Future
import scala.concurrent.duration.Duration
@@ -88,7 +87,7 @@ object CassandraProjection {
def groupedWithin[Offset, Envelope](
projectionId: ProjectionId,
sourceProvider: SourceProvider[Offset, Envelope],
- handler: () => Handler[immutable.Seq[Envelope]]):
GroupedProjection[Offset, Envelope] =
+ handler: () => Handler[Seq[Envelope]]): GroupedProjection[Offset,
Envelope] =
new CassandraProjectionImpl(
projectionId,
sourceProvider,
diff --git
a/core-test/src/test/scala/org/apache/pekko/projection/internal/OffsetSerializationSpec.scala
b/core-test/src/test/scala/org/apache/pekko/projection/internal/OffsetSerializationSpec.scala
index 74c490a8..32a9b4f6 100644
---
a/core-test/src/test/scala/org/apache/pekko/projection/internal/OffsetSerializationSpec.scala
+++
b/core-test/src/test/scala/org/apache/pekko/projection/internal/OffsetSerializationSpec.scala
@@ -17,8 +17,6 @@ import java.nio.charset.StandardCharsets
import java.util.Base64
import java.util.UUID
-import scala.collection.immutable
-
import org.apache.pekko
import pekko.actor.ExtendedActorSystem
import pekko.actor.testkit.typed.scaladsl.LogCapturing
@@ -140,7 +138,7 @@ class OffsetSerializationSpec
}
val storageRepresentation = MultipleOffsets(
- immutable.Seq(SingleOffset(ProjectionId(id.name,
surrogateProjectionKey), LongManifest, "1", mergeable = true)))
+ Seq(SingleOffset(ProjectionId(id.name, surrogateProjectionKey),
LongManifest, "1", mergeable = true)))
actualRep shouldBe storageRepresentation
@@ -155,7 +153,7 @@ class OffsetSerializationSpec
val mergeableOffset =
MergeableOffset(Map(surrogateProjectionKey1 -> 1L,
surrogateProjectionKey2 -> 2L))
val storageRepresentation = MultipleOffsets(
- immutable.Seq(
+ Seq(
SingleOffset(ProjectionId(projectionName, surrogateProjectionKey1),
LongManifest, "1", mergeable = true),
SingleOffset(ProjectionId(projectionName, surrogateProjectionKey2),
LongManifest, "2", mergeable = true)))
diff --git
a/core-test/src/test/scala/org/apache/pekko/projection/internal/metrics/tools/InternalProjectionStateMetricsSpec.scala
b/core-test/src/test/scala/org/apache/pekko/projection/internal/metrics/tools/InternalProjectionStateMetricsSpec.scala
index 674a0974..88609e7d 100644
---
a/core-test/src/test/scala/org/apache/pekko/projection/internal/metrics/tools/InternalProjectionStateMetricsSpec.scala
+++
b/core-test/src/test/scala/org/apache/pekko/projection/internal/metrics/tools/InternalProjectionStateMetricsSpec.scala
@@ -15,7 +15,6 @@ package org.apache.pekko.projection.internal.metrics.tools
import java.util.UUID
-import scala.collection.immutable
import scala.concurrent.Await
import scala.concurrent.ExecutionContext
import scala.concurrent.Future
@@ -161,8 +160,8 @@ object InternalProjectionStateMetricsSpec {
}
case groupedHandlerStrategy: GroupedHandlerStrategy[Envelope]
@unchecked => {
val adaptedHandler = () =>
- new Handler[immutable.Seq[Envelope]] {
- override def process(envelopes: immutable.Seq[Envelope]):
Future[Done] =
+ new Handler[Seq[Envelope]] {
+ override def process(envelopes: Seq[Envelope]): Future[Done] =
groupedHandlerStrategy.handlerFactory().process(envelopes).flatMap { _ =>
offsetStore.saveOffset(projectionId, envelopes.last.offset)
}
diff --git
a/core-test/src/test/scala/org/apache/pekko/projection/internal/metrics/tools/TestHandlers.scala
b/core-test/src/test/scala/org/apache/pekko/projection/internal/metrics/tools/TestHandlers.scala
index 1df7150a..c51ea84b 100644
---
a/core-test/src/test/scala/org/apache/pekko/projection/internal/metrics/tools/TestHandlers.scala
+++
b/core-test/src/test/scala/org/apache/pekko/projection/internal/metrics/tools/TestHandlers.scala
@@ -13,7 +13,6 @@
package org.apache.pekko.projection.internal.metrics.tools
-import scala.collection.immutable
import scala.concurrent.Future
import org.apache.pekko
@@ -71,13 +70,13 @@ object TestHandlers {
* trigger an error and then be removed from the
stack. To fail an item multiple times
* add its offset repeatedly. Uses `Int` instead of
`Long` for convenience.
*/
- def groupedWithErrors(erroredOffsets: Int*): () =>
Handler[immutable.Seq[Envelope]] = {
+ def groupedWithErrors(erroredOffsets: Int*): () => Handler[Seq[Envelope]] = {
var nextProcessStrategy = ProcessStrategy(erroredOffsets.map {
_.toLong
}.toList)
() =>
- new Handler[immutable.Seq[Envelope]] {
- override def process(envelopes: immutable.Seq[Envelope]): Future[Done]
= {
+ new Handler[Seq[Envelope]] {
+ override def process(envelopes: Seq[Envelope]): Future[Done] = {
nextProcessStrategy match {
case SomeFailures(nextFail :: tail)
if envelopes
diff --git
a/core/src/main/scala/org/apache/pekko/projection/internal/HandlerAdapter.scala
b/core/src/main/scala/org/apache/pekko/projection/internal/HandlerAdapter.scala
index f7ab5496..3d650f63 100644
---
a/core/src/main/scala/org/apache/pekko/projection/internal/HandlerAdapter.scala
+++
b/core/src/main/scala/org/apache/pekko/projection/internal/HandlerAdapter.scala
@@ -13,7 +13,6 @@
package org.apache.pekko.projection.internal
-import scala.collection.immutable
import scala.concurrent.Future
import scala.jdk.CollectionConverters._
import scala.jdk.FutureConverters._
@@ -55,13 +54,13 @@ import pekko.projection.scaladsl
}
/**
- * INTERNAL API: Adapter from `javadsl.Handler[java.util.List[Envelope]]` to
`scaladsl.Handler[immutable.Seq[Envelope]]`
+ * INTERNAL API: Adapter from `javadsl.Handler[java.util.List[Envelope]]` to
`scaladsl.Handler[Seq[Envelope]]`
*/
@InternalApi private[projection] class GroupedHandlerAdapter[Envelope](
delegate: javadsl.Handler[java.util.List[Envelope]])
- extends scaladsl.Handler[immutable.Seq[Envelope]] {
+ extends scaladsl.Handler[Seq[Envelope]] {
- override def process(envelopes: immutable.Seq[Envelope]): Future[Done] = {
+ override def process(envelopes: Seq[Envelope]): Future[Done] = {
delegate.process(envelopes.asJava).asScala
}
diff --git
a/core/src/main/scala/org/apache/pekko/projection/internal/InternalProjectionState.scala
b/core/src/main/scala/org/apache/pekko/projection/internal/InternalProjectionState.scala
index e4a642e4..badea35a 100644
---
a/core/src/main/scala/org/apache/pekko/projection/internal/InternalProjectionState.scala
+++
b/core/src/main/scala/org/apache/pekko/projection/internal/InternalProjectionState.scala
@@ -13,7 +13,6 @@
package org.apache.pekko.projection.internal
-import scala.collection.immutable
import scala.concurrent.ExecutionContext
import scala.concurrent.Future
import scala.concurrent.Promise
@@ -89,7 +88,7 @@ private[projection] abstract class
InternalProjectionState[Offset, Envelope](
protected def saveOffsetsAndReport(
projectionId: ProjectionId,
- batch: immutable.Seq[ProjectionContextImpl[Offset, Envelope]]):
Future[Done] = {
+ batch: Seq[ProjectionContextImpl[Offset, Envelope]]): Future[Done] = {
// The batch contains multiple projections contexts. Each of these
contexts may represent
// a single envelope or a group of envelopes. The size of the batch and
the size of the
@@ -103,15 +102,15 @@ private[projection] abstract class
InternalProjectionState[Offset, Envelope](
/**
* A convenience method to serialize asynchronous operations to occur one
after another is complete
*/
- private def serialize(batches: Map[String,
immutable.Seq[ProjectionContextImpl[Offset, Envelope]]])(
- op: (String, immutable.Seq[ProjectionContextImpl[Offset, Envelope]]) =>
Future[Done]): Future[Done] = {
+ private def serialize(batches: Map[String, Seq[ProjectionContextImpl[Offset,
Envelope]]])(
+ op: (String, Seq[ProjectionContextImpl[Offset, Envelope]]) =>
Future[Done]): Future[Done] = {
val logProgressEvery: Int = 5
val size = batches.size
logger.debug("Processing [{}] partitioned batches serially", size)
def loop(
- remaining: List[(String, immutable.Seq[ProjectionContextImpl[Offset,
Envelope]])],
+ remaining: List[(String, Seq[ProjectionContextImpl[Offset,
Envelope]])],
n: Int): Future[Done] = {
remaining match {
case Nil => Future.successful(Done)
@@ -247,11 +246,11 @@ private[projection] abstract class
InternalProjectionState[Offset, Envelope](
HandlerRecoveryImpl[Offset, Envelope](projectionId, recoveryStrategy,
logger, statusObserver, telemetry)
def processGrouped(
- handler: Handler[immutable.Seq[Envelope]],
+ handler: Handler[Seq[Envelope]],
handlerRecovery: HandlerRecoveryImpl[Offset, Envelope],
- envelopesAndOffsets: immutable.Seq[ProjectionContextImpl[Offset,
Envelope]]): Future[Done] = {
+ envelopesAndOffsets: Seq[ProjectionContextImpl[Offset, Envelope]]):
Future[Done] = {
- def processEnvelopes(partitioned:
immutable.Seq[ProjectionContextImpl[Offset, Envelope]]): Future[Done] = {
+ def processEnvelopes(partitioned: Seq[ProjectionContextImpl[Offset,
Envelope]]): Future[Done] = {
val first = partitioned.head
val firstOffset = first.offset
val lastOffset = partitioned.last.offset
diff --git
a/core/src/main/scala/org/apache/pekko/projection/internal/OffsetSerialization.scala
b/core/src/main/scala/org/apache/pekko/projection/internal/OffsetSerialization.scala
index 6dbda1d2..f87de96f 100644
---
a/core/src/main/scala/org/apache/pekko/projection/internal/OffsetSerialization.scala
+++
b/core/src/main/scala/org/apache/pekko/projection/internal/OffsetSerialization.scala
@@ -16,8 +16,6 @@ package org.apache.pekko.projection.internal
import java.util.Base64
import java.util.UUID
-import scala.collection.immutable
-
import org.apache.pekko
import pekko.actor.typed.ActorSystem
import pekko.annotation.InternalApi
@@ -34,7 +32,7 @@ import pekko.serialization.Serializers
sealed trait StorageRepresentation
final case class SingleOffset(id: ProjectionId, manifest: String, offsetStr:
String, mergeable: Boolean = false)
extends StorageRepresentation
- final case class MultipleOffsets(reps: immutable.Seq[SingleOffset]) extends
StorageRepresentation
+ final case class MultipleOffsets(reps: Seq[SingleOffset]) extends
StorageRepresentation
final val StringManifest = "STR"
final val LongManifest = "LNG"
diff --git
a/core/src/main/scala/org/apache/pekko/projection/internal/OffsetStrategy.scala
b/core/src/main/scala/org/apache/pekko/projection/internal/OffsetStrategy.scala
index b6d26937..b9d034e5 100644
---
a/core/src/main/scala/org/apache/pekko/projection/internal/OffsetStrategy.scala
+++
b/core/src/main/scala/org/apache/pekko/projection/internal/OffsetStrategy.scala
@@ -13,7 +13,6 @@
package org.apache.pekko.projection.internal
-import scala.collection.immutable
import scala.concurrent.duration.FiniteDuration
import org.apache.pekko
@@ -128,10 +127,10 @@ private[projection] final case class
SingleHandlerStrategy[Envelope](handlerFact
*/
@InternalApi
private[projection] final case class GroupedHandlerStrategy[Envelope](
- handlerFactory: () => Handler[immutable.Seq[Envelope]],
+ handlerFactory: () => Handler[Seq[Envelope]],
afterEnvelopes: Option[Int] = None,
orAfterDuration: Option[FiniteDuration] = None)
- extends FunctionHandlerStrategy[immutable.Seq[Envelope]](handlerFactory)
+ extends FunctionHandlerStrategy[Seq[Envelope]](handlerFactory)
/**
* INTERNAL API
diff --git
a/core/src/main/scala/org/apache/pekko/projection/internal/Telemetry.scala
b/core/src/main/scala/org/apache/pekko/projection/internal/Telemetry.scala
index e3542642..b610c92c 100644
--- a/core/src/main/scala/org/apache/pekko/projection/internal/Telemetry.scala
+++ b/core/src/main/scala/org/apache/pekko/projection/internal/Telemetry.scala
@@ -15,7 +15,6 @@ package org.apache.pekko.projection.internal
import java.util
-import scala.collection.immutable
import scala.jdk.CollectionConverters._
import org.apache.pekko
@@ -122,7 +121,7 @@ trait Telemetry {
dynamicAccess
.createInstanceFor[Telemetry](
fqcn,
- immutable.Seq((classOf[ProjectionId], projectionId),
(classOf[ActorSystem[?]], system)))
+ Seq((classOf[ProjectionId], projectionId), (classOf[ActorSystem[?]],
system)))
.get
}
}
diff --git a/docs/src/main/paradox/cassandra.md
b/docs/src/main/paradox/cassandra.md
index 608c92f3..58c0d5cb 100644
--- a/docs/src/main/paradox/cassandra.md
+++ b/docs/src/main/paradox/cassandra.md
@@ -82,7 +82,7 @@ The envelopes are grouped within a time window, or limited by
a number of envelo
This window can be defined with `withGroup` of the returned
`GroupedProjection`. The default settings for
the window is defined in configuration section `pekko.projection.grouped`.
-When using `groupedWithin` the handler is a
@scala[`Handler[immutable.Seq[EventEnvelope[ShoppingCart.Event]]]`]@java[`Handler<List<EventEnvelope<ShoppingCart.Event>>>`].
+When using `groupedWithin` the handler is a
@scala[`Handler[Seq[EventEnvelope[ShoppingCart.Event]]]`]@java[`Handler<List<EventEnvelope<ShoppingCart.Event>>>`].
The @ref:[`GroupedShoppingCartHandler` is shown below](#grouped-handler).
It stores the offset in Cassandra immediately after the `handler` has
processed the envelopes, but that
diff --git a/docs/src/main/paradox/jdbc.md b/docs/src/main/paradox/jdbc.md
index c504a660..fd968f7f 100644
--- a/docs/src/main/paradox/jdbc.md
+++ b/docs/src/main/paradox/jdbc.md
@@ -125,7 +125,7 @@ The envelopes are grouped within a time window, or limited
by a number of envelo
This window can be defined with `withGroup` of the returned
`GroupedProjection`. The default settings for
the window is defined in configuration section `pekko.projection.grouped`.
-When using `groupedWithin` the handler is a
@scala[`JdbcHandler[immutable.Seq[EventEnvelope[ShoppingCart.Event]]]`]@java[`JdbcHandler<List<EventEnvelope<ShoppingCart.Event>>>`].
+When using `groupedWithin` the handler is a
@scala[`JdbcHandler[Seq[EventEnvelope[ShoppingCart.Event]]]`]@java[`JdbcHandler<List<EventEnvelope<ShoppingCart.Event>>>`].
The @ref:[`GroupedShoppingCartHandler` is shown below](#grouped-handler).
The offset is stored in the same transaction used for the user defined
`handler`, which means exactly-once
diff --git a/docs/src/main/paradox/r2dbc.md b/docs/src/main/paradox/r2dbc.md
index 34a238d4..5211fa7e 100644
--- a/docs/src/main/paradox/r2dbc.md
+++ b/docs/src/main/paradox/r2dbc.md
@@ -160,7 +160,7 @@ The envelopes are grouped within a time window, or limited
by a number of envelo
This window can be defined with `withGroup` of the returned
`GroupedProjection`. The default settings for
the window is defined in configuration section `pekko.projection.grouped`.
-When using `groupedWithin` the handler is a
@scala[`R2dbcHandler[immutable.Seq[EventEnvelope[ShoppingCart.Event]]]`]@java[`R2dbcHandler<List<EventEnvelope<ShoppingCart.Event>>>`].
+When using `groupedWithin` the handler is a
@scala[`R2dbcHandler[Seq[EventEnvelope[ShoppingCart.Event]]]`]@java[`R2dbcHandler<List<EventEnvelope<ShoppingCart.Event>>>`].
The @ref:[`GroupedShoppingCartHandler` is shown below](#grouped-handler).
The offset is stored in the same transaction used for the user defined
`handler`, which means exactly-once
diff --git a/docs/src/main/paradox/slick.md b/docs/src/main/paradox/slick.md
index 420f81d9..821e4d2e 100644
--- a/docs/src/main/paradox/slick.md
+++ b/docs/src/main/paradox/slick.md
@@ -74,7 +74,7 @@ The envelopes are grouped within a time window, or limited by
a number of envelo
This window can be defined with `withGroup` of the returned
`GroupedProjection`. The default settings for
the window is defined in configuration section `pekko.projection.grouped`.
-When using `groupedWithin` the handler is a
`SlickHandler[immutable.Seq[EventEnvelope[ShoppingCart.Event]]]`.
+When using `groupedWithin` the handler is a
`SlickHandler[Seq[EventEnvelope[ShoppingCart.Event]]]`.
The @ref:[`GroupedShoppingCartHandler` is shown below](#grouped-handler).
The offset is stored in the same transaction as the `DBIO` returned from the
`handler`, which means exactly-once
diff --git
a/durable-state/src/main/scala/org/apache/pekko/projection/state/scaladsl/DurableStateSourceProvider.scala
b/durable-state/src/main/scala/org/apache/pekko/projection/state/scaladsl/DurableStateSourceProvider.scala
index ab3ea964..62828bf4 100644
---
a/durable-state/src/main/scala/org/apache/pekko/projection/state/scaladsl/DurableStateSourceProvider.scala
+++
b/durable-state/src/main/scala/org/apache/pekko/projection/state/scaladsl/DurableStateSourceProvider.scala
@@ -13,7 +13,6 @@
package org.apache.pekko.projection.state.scaladsl
-import scala.collection.immutable
import scala.concurrent.ExecutionContext
import scala.concurrent.Future
@@ -112,7 +111,7 @@ object DurableStateSourceProvider {
def sliceRanges(
system: ActorSystem[?],
durableStateStoreQueryPluginId: String,
- numberOfRanges: Int): immutable.Seq[Range] =
+ numberOfRanges: Int): Seq[Range] =
DurableStateStoreRegistry(system)
.durableStateStoreFor[DurableStateStoreBySliceQuery[Any]](durableStateStoreQueryPluginId)
.sliceRanges(numberOfRanges)
diff --git
a/eventsourced/src/main/scala/org/apache/pekko/projection/eventsourced/scaladsl/EventSourcedProvider.scala
b/eventsourced/src/main/scala/org/apache/pekko/projection/eventsourced/scaladsl/EventSourcedProvider.scala
index 3668b051..3778af21 100644
---
a/eventsourced/src/main/scala/org/apache/pekko/projection/eventsourced/scaladsl/EventSourcedProvider.scala
+++
b/eventsourced/src/main/scala/org/apache/pekko/projection/eventsourced/scaladsl/EventSourcedProvider.scala
@@ -15,7 +15,6 @@ package org.apache.pekko.projection.eventsourced.scaladsl
import java.time.Instant
-import scala.collection.immutable
import scala.concurrent.ExecutionContext
import scala.concurrent.Future
import com.typesafe.config.Config
@@ -141,7 +140,7 @@ object EventSourcedProvider {
.readJournalFor[EventsBySliceQuery](readJournalPluginId,
readJournalConfig)
.sliceForPersistenceId(persistenceId)
- def sliceRanges(system: ActorSystem[?], readJournalPluginId: String,
numberOfRanges: Int): immutable.Seq[Range] =
+ def sliceRanges(system: ActorSystem[?], readJournalPluginId: String,
numberOfRanges: Int): Seq[Range] =
PersistenceQuery(system).readJournalFor[EventsBySliceQuery](readJournalPluginId).sliceRanges(numberOfRanges)
/** @since 2.0.0 */
@@ -149,7 +148,7 @@ object EventSourcedProvider {
system: ActorSystem[?],
readJournalPluginId: String,
readJournalConfig: Config,
- numberOfRanges: Int): immutable.Seq[Range] =
+ numberOfRanges: Int): Seq[Range] =
PersistenceQuery(system)
.readJournalFor[EventsBySliceQuery](readJournalPluginId,
readJournalConfig)
.sliceRanges(numberOfRanges)
diff --git
a/eventsourced/src/test/scala/org/apache/pekko/projection/eventsourced/scaldsl/EventSourcedProviderSpec.scala
b/eventsourced/src/test/scala/org/apache/pekko/projection/eventsourced/scaldsl/EventSourcedProviderSpec.scala
index c5f2e16e..8b87ee0b 100644
---
a/eventsourced/src/test/scala/org/apache/pekko/projection/eventsourced/scaldsl/EventSourcedProviderSpec.scala
+++
b/eventsourced/src/test/scala/org/apache/pekko/projection/eventsourced/scaldsl/EventSourcedProviderSpec.scala
@@ -17,7 +17,6 @@
package org.apache.pekko.projection.eventsourced.scaldsl
-import scala.collection.immutable.Seq
import scala.concurrent.Future
import com.typesafe.config.ConfigFactory
import org.apache.pekko
diff --git a/examples/src/test/scala/docs/jdbc/JdbcProjectionDocExample.scala
b/examples/src/test/scala/docs/jdbc/JdbcProjectionDocExample.scala
index 87bf7498..eeaecfad 100644
--- a/examples/src/test/scala/docs/jdbc/JdbcProjectionDocExample.scala
+++ b/examples/src/test/scala/docs/jdbc/JdbcProjectionDocExample.scala
@@ -99,15 +99,13 @@ object JdbcProjectionDocExample {
// #handler
// #grouped-handler
- import scala.collection.immutable
-
class GroupedShoppingCartHandler(repository: OrderRepository)
- extends JdbcHandler[immutable.Seq[EventEnvelope[ShoppingCart.Event]],
PlainJdbcSession] {
+ extends JdbcHandler[Seq[EventEnvelope[ShoppingCart.Event]],
PlainJdbcSession] {
private val logger = LoggerFactory.getLogger(getClass)
override def process(
session: PlainJdbcSession,
- envelopes: immutable.Seq[EventEnvelope[ShoppingCart.Event]]): Unit = {
+ envelopes: Seq[EventEnvelope[ShoppingCart.Event]]): Unit = {
// save all events in DB
envelopes.map(_.event).foreach {
diff --git a/examples/src/test/scala/docs/slick/SlickProjectionDocExample.scala
b/examples/src/test/scala/docs/slick/SlickProjectionDocExample.scala
index 4006c77e..9ae1f84c 100644
--- a/examples/src/test/scala/docs/slick/SlickProjectionDocExample.scala
+++ b/examples/src/test/scala/docs/slick/SlickProjectionDocExample.scala
@@ -94,13 +94,11 @@ class SlickProjectionDocExample {
// #handler
// #grouped-handler
- import scala.collection.immutable
-
class GroupedShoppingCartHandler(repository: OrderRepository)(implicit ec:
ExecutionContext)
- extends SlickHandler[immutable.Seq[EventEnvelope[ShoppingCart.Event]]] {
+ extends SlickHandler[Seq[EventEnvelope[ShoppingCart.Event]]] {
private val logger = LoggerFactory.getLogger(getClass)
- override def process(envelopes:
immutable.Seq[EventEnvelope[ShoppingCart.Event]]): DBIO[Done] = {
+ override def process(envelopes: Seq[EventEnvelope[ShoppingCart.Event]]):
DBIO[Done] = {
val dbios = envelopes.map(_.event).map {
case ShoppingCart.CheckedOut(cartId, time) =>
logger.info(s"Shopping cart $cartId was checked out at $time")
diff --git
a/grpc-test/src/test/scala/org/apache/pekko/projection/grpc/internal/EventProducerServiceSpec.scala
b/grpc-test/src/test/scala/org/apache/pekko/projection/grpc/internal/EventProducerServiceSpec.scala
index d15518c9..0740508c 100644
---
a/grpc-test/src/test/scala/org/apache/pekko/projection/grpc/internal/EventProducerServiceSpec.scala
+++
b/grpc-test/src/test/scala/org/apache/pekko/projection/grpc/internal/EventProducerServiceSpec.scala
@@ -54,7 +54,6 @@ import org.scalatest.wordspec.AnyWordSpecLike
import java.time.Instant
import java.util.concurrent.ConcurrentHashMap
-import scala.collection.immutable
import scala.concurrent.Future
import scala.concurrent.Promise
@@ -92,7 +91,7 @@ object EventProducerServiceSpec {
override def sliceForPersistenceId(persistenceId: String): Int =
persistenceExt.sliceForPersistenceId(persistenceId)
- override def sliceRanges(numberOfRanges: Int): immutable.Seq[Range] =
+ override def sliceRanges(numberOfRanges: Int): Seq[Range] =
persistenceExt.sliceRanges(numberOfRanges)
}
}
diff --git
a/grpc/src/main/scala/org/apache/pekko/projection/grpc/consumer/ConsumerFilter.scala
b/grpc/src/main/scala/org/apache/pekko/projection/grpc/consumer/ConsumerFilter.scala
index c588e114..36409959 100644
---
a/grpc/src/main/scala/org/apache/pekko/projection/grpc/consumer/ConsumerFilter.scala
+++
b/grpc/src/main/scala/org/apache/pekko/projection/grpc/consumer/ConsumerFilter.scala
@@ -17,7 +17,6 @@ import java.util.{ List => JList }
import java.util.{ Set => JSet }
import scala.annotation.tailrec
-import scala.collection.immutable
import scala.concurrent.duration.FiniteDuration
import org.apache.pekko
@@ -57,7 +56,7 @@ object ConsumerFilter extends ExtensionId[ConsumerFilter] {
*/
@InternalApi private[pekko] final case class Subscribe(
streamId: String,
- initCriteria: immutable.Seq[FilterCriteria],
+ initCriteria: Seq[FilterCriteria],
subscriber: ActorRef[SubscriberCommand])
extends Command
@@ -70,7 +69,7 @@ object ConsumerFilter extends ExtensionId[ConsumerFilter] {
* If no matching include criteria the event is discarded.
* If matching include criteria the event is emitted.
*/
- final case class UpdateFilter(streamId: String, criteria:
immutable.Seq[FilterCriteria]) extends SubscriberCommand {
+ final case class UpdateFilter(streamId: String, criteria:
Seq[FilterCriteria]) extends SubscriberCommand {
/** Java API */
def this(streamId: String, criteria: JList[FilterCriteria]) =
@@ -79,7 +78,7 @@ object ConsumerFilter extends ExtensionId[ConsumerFilter] {
final case class GetFilter(streamId: String, replyTo:
ActorRef[CurrentFilter]) extends Command
- final case class CurrentFilter(streamId: String, criteria:
immutable.Seq[FilterCriteria]) {
+ final case class CurrentFilter(streamId: String, criteria:
Seq[FilterCriteria]) {
/** Java API */
def getCriteria(): JList[FilterCriteria] =
@@ -295,8 +294,8 @@ object ConsumerFilter extends ExtensionId[ConsumerFilter] {
* INTERNAL API
*/
@InternalApi private[pekko] def mergeFilter(
- currentFilter: immutable.Seq[FilterCriteria],
- update: immutable.Seq[FilterCriteria]): immutable.Seq[FilterCriteria] = {
+ currentFilter: Seq[FilterCriteria],
+ update: Seq[FilterCriteria]): Seq[FilterCriteria] = {
val both = currentFilter ++ update
@@ -376,8 +375,8 @@ object ConsumerFilter extends ExtensionId[ConsumerFilter] {
* INTERNAL API
*/
@InternalApi private[pekko] def createDiff(
- a: immutable.Seq[FilterCriteria],
- b: immutable.Seq[FilterCriteria]): immutable.Seq[FilterCriteria] = {
+ a: Seq[FilterCriteria],
+ b: Seq[FilterCriteria]): Seq[FilterCriteria] = {
require(!hasRemoveCriteria(a), "Unexpected RemoveCriteria in a when
creating diff, use mergeFilter first.")
require(!hasRemoveCriteria(b), "Unexpected RemoveCriteria in b when
creating diff, use mergeFilter first.")
@@ -461,7 +460,7 @@ object ConsumerFilter extends ExtensionId[ConsumerFilter] {
}
/** INTERNAL API */
- @InternalApi private[pekko] def includeEntityOffsets(filter:
immutable.Seq[FilterCriteria]): Set[EntityIdOffset] = {
+ @InternalApi private[pekko] def includeEntityOffsets(filter:
Seq[FilterCriteria]): Set[EntityIdOffset] = {
filter.flatMap {
case inc: IncludeEntityIds => inc.entityOffsets
case _ => Set.empty[EntityIdOffset]
@@ -469,7 +468,7 @@ object ConsumerFilter extends ExtensionId[ConsumerFilter] {
}
/** INTERNAL API */
- @InternalApi private[pekko] def excludeTags(filter:
immutable.Seq[FilterCriteria]): Set[String] = {
+ @InternalApi private[pekko] def excludeTags(filter: Seq[FilterCriteria]):
Set[String] = {
filter.flatMap {
case exl: ExcludeTags => exl.tags
case _ => Set.empty[String]
@@ -477,7 +476,7 @@ object ConsumerFilter extends ExtensionId[ConsumerFilter] {
}
/** INTERNAL API */
- @InternalApi private[pekko] def includeTags(filter:
immutable.Seq[FilterCriteria]): Set[String] = {
+ @InternalApi private[pekko] def includeTags(filter: Seq[FilterCriteria]):
Set[String] = {
filter.flatMap {
case incl: IncludeTags => incl.tags
case _ => Set.empty[String]
@@ -485,7 +484,7 @@ object ConsumerFilter extends ExtensionId[ConsumerFilter] {
}
/** INTERNAL API */
- @InternalApi private[pekko] def excludeEntityIds(filter:
immutable.Seq[FilterCriteria]): Set[String] = {
+ @InternalApi private[pekko] def excludeEntityIds(filter:
Seq[FilterCriteria]): Set[String] = {
filter.flatMap {
case exl: ExcludeEntityIds => exl.entityIds
case _ => Set.empty[String]
@@ -493,7 +492,7 @@ object ConsumerFilter extends ExtensionId[ConsumerFilter] {
}
/** INTERNAL API */
- @InternalApi private[pekko] def excludeRegexEntityIds(filter:
immutable.Seq[FilterCriteria]): Set[String] = {
+ @InternalApi private[pekko] def excludeRegexEntityIds(filter:
Seq[FilterCriteria]): Set[String] = {
filter.flatMap {
case rxp: ExcludeRegexEntityIds => rxp.matching
case _ => Set.empty[String]
@@ -501,7 +500,7 @@ object ConsumerFilter extends ExtensionId[ConsumerFilter] {
}
/** INTERNAL API */
- @InternalApi private[pekko] def includeRegexEntityIds(filter:
immutable.Seq[FilterCriteria]): Set[String] = {
+ @InternalApi private[pekko] def includeRegexEntityIds(filter:
Seq[FilterCriteria]): Set[String] = {
filter.flatMap {
case rxp: IncludeRegexEntityIds => rxp.matching
case _ => Set.empty[String]
@@ -509,7 +508,7 @@ object ConsumerFilter extends ExtensionId[ConsumerFilter] {
}
/** INTERNAL API */
- @InternalApi private[pekko] def hasRemoveCriteria(filter:
immutable.Seq[FilterCriteria]): Boolean =
+ @InternalApi private[pekko] def hasRemoveCriteria(filter:
Seq[FilterCriteria]): Boolean =
filter.exists(_.isInstanceOf[RemoveCriteria])
/** INTERNAL API */
diff --git
a/grpc/src/main/scala/org/apache/pekko/projection/grpc/consumer/scaladsl/GrpcReadJournal.scala
b/grpc/src/main/scala/org/apache/pekko/projection/grpc/consumer/scaladsl/GrpcReadJournal.scala
index 58834280..236af4d1 100644
---
a/grpc/src/main/scala/org/apache/pekko/projection/grpc/consumer/scaladsl/GrpcReadJournal.scala
+++
b/grpc/src/main/scala/org/apache/pekko/projection/grpc/consumer/scaladsl/GrpcReadJournal.scala
@@ -16,7 +16,6 @@ package org.apache.pekko.projection.grpc.consumer.scaladsl
import java.time.Instant
import java.util.concurrent.TimeUnit
-import scala.collection.immutable
import scala.concurrent.ExecutionContext
import scala.concurrent.Future
@@ -101,7 +100,7 @@ object GrpcReadJournal {
* Note that the `protobufDescriptors` is a list of the `javaDescriptor` for
the used protobuf messages. It is
* defined in the ScalaPB generated `Proto` companion object.
*/
- def apply(protobufDescriptors: immutable.Seq[Descriptors.FileDescriptor])(
+ def apply(protobufDescriptors: Seq[Descriptors.FileDescriptor])(
implicit system: ClassicActorSystemProvider): GrpcReadJournal =
apply(
GrpcQuerySettings(system),
@@ -118,7 +117,7 @@ object GrpcReadJournal {
def apply(
settings: GrpcQuerySettings,
clientSettings: GrpcClientSettings,
- protobufDescriptors: immutable.Seq[Descriptors.FileDescriptor])(
+ protobufDescriptors: Seq[Descriptors.FileDescriptor])(
implicit system: ClassicActorSystemProvider): GrpcReadJournal =
apply(settings, clientSettings, protobufDescriptors,
ProtoAnySerialization.Prefer.Scala)
@@ -128,7 +127,7 @@ object GrpcReadJournal {
@InternalApi private[pekko] def apply(
settings: GrpcQuerySettings,
clientSettings: GrpcClientSettings,
- protobufDescriptors: immutable.Seq[Descriptors.FileDescriptor],
+ protobufDescriptors: Seq[Descriptors.FileDescriptor],
protobufPrefer: ProtoAnySerialization.Prefer)(implicit system:
ClassicActorSystemProvider): GrpcReadJournal = {
// FIXME issue #702 This probably means that one GrpcReadJournal instance
is created for each Projection instance,
@@ -220,7 +219,7 @@ final class GrpcReadJournal private (
override def sliceForPersistenceId(persistenceId: String): Int =
persistenceExt.sliceForPersistenceId(persistenceId)
- override def sliceRanges(numberOfRanges: Int): immutable.Seq[Range] =
+ override def sliceRanges(numberOfRanges: Int): Seq[Range] =
persistenceExt.sliceRanges(numberOfRanges)
/**
@@ -296,7 +295,7 @@ final class GrpcReadJournal private (
throw new IllegalArgumentException(s"Expected TimestampOffset or
NoOffset, but got [$offset]")
}
- def inReqSource(initCriteria:
immutable.Seq[ConsumerFilter.FilterCriteria]): Source[StreamIn, NotUsed] =
+ def inReqSource(initCriteria: Seq[ConsumerFilter.FilterCriteria]):
Source[StreamIn, NotUsed] =
Source
.actorRef[ConsumerFilter.SubscriberCommand](
completionMatcher = PartialFunction.empty,
@@ -390,7 +389,7 @@ final class GrpcReadJournal private (
}
}
- private def toProtoFilterCriteria(criteria:
immutable.Seq[ConsumerFilter.FilterCriteria]): Seq[FilterCriteria] = {
+ private def toProtoFilterCriteria(criteria:
Seq[ConsumerFilter.FilterCriteria]): Seq[FilterCriteria] = {
criteria.map {
case ConsumerFilter.ExcludeTags(tags) =>
FilterCriteria(FilterCriteria.Message.ExcludeTags(ExcludeTags(tags.toVector)))
diff --git
a/grpc/src/main/scala/org/apache/pekko/projection/grpc/internal/ConsumerFilterRegistry.scala
b/grpc/src/main/scala/org/apache/pekko/projection/grpc/internal/ConsumerFilterRegistry.scala
index 14eb829a..3994c777 100644
---
a/grpc/src/main/scala/org/apache/pekko/projection/grpc/internal/ConsumerFilterRegistry.scala
+++
b/grpc/src/main/scala/org/apache/pekko/projection/grpc/internal/ConsumerFilterRegistry.scala
@@ -16,7 +16,6 @@ package org.apache.pekko.projection.grpc.internal
import java.net.URLEncoder
import java.nio.charset.StandardCharsets
-import scala.collection.immutable
import scala.concurrent.duration._
import scala.util.Failure
import scala.util.Success
@@ -39,7 +38,7 @@ import pekko.util.Timeout
sealed trait InternalCommand extends Command
- final case class FilterUpdated(streamId: String, criteria:
immutable.Seq[FilterCriteria]) extends InternalCommand
+ final case class FilterUpdated(streamId: String, criteria:
Seq[FilterCriteria]) extends InternalCommand
private final case class SubscriberTerminated(subscriber: Subscriber)
extends InternalCommand
@@ -64,7 +63,7 @@ import pekko.util.Timeout
import ConsumerFilterRegistry._
private def behavior(
- subscribers: Map[Subscriber, immutable.Seq[FilterCriteria]],
+ subscribers: Map[Subscriber, Seq[FilterCriteria]],
stores: Map[String, ActorRef[ConsumerFilterStore.Command]]):
Behavior[Command] = {
def getOrCreateStore(streamId: String):
ActorRef[ConsumerFilterStore.Command] = {
@@ -79,7 +78,7 @@ import pekko.util.Timeout
def publishUpdatedFilterToSubscribers(
streamId: String,
- filter: immutable.Seq[FilterCriteria]): Map[Subscriber,
immutable.Seq[FilterCriteria]] = {
+ filter: Seq[FilterCriteria]): Map[Subscriber, Seq[FilterCriteria]] = {
subscribers.map {
case (sub, subFilter) =>
if (sub.streamId == streamId) {
diff --git
a/grpc/src/main/scala/org/apache/pekko/projection/grpc/internal/ConsumerFilterStore.scala
b/grpc/src/main/scala/org/apache/pekko/projection/grpc/internal/ConsumerFilterStore.scala
index 641069eb..ffbb18ef 100644
---
a/grpc/src/main/scala/org/apache/pekko/projection/grpc/internal/ConsumerFilterStore.scala
+++
b/grpc/src/main/scala/org/apache/pekko/projection/grpc/internal/ConsumerFilterStore.scala
@@ -16,7 +16,6 @@ package org.apache.pekko.projection.grpc.internal
import java.util.ConcurrentModificationException
import java.util.concurrent.ConcurrentHashMap
-import scala.collection.immutable
import scala.util.Failure
import scala.util.Success
@@ -54,7 +53,7 @@ import org.slf4j.LoggerFactory
@InternalApi private[pekko] object ConsumerFilterStore {
sealed trait Command
- final case class UpdateFilter(criteria: immutable.Seq[FilterCriteria])
extends Command
+ final case class UpdateFilter(criteria: Seq[FilterCriteria]) extends Command
final case class GetFilter(replyTo: ActorRef[ConsumerFilter.CurrentFilter])
extends Command
@@ -117,7 +116,7 @@ import org.slf4j.LoggerFactory
}
private class StoreExt extends Extension {
- val filtersByStreamId = new ConcurrentHashMap[String,
immutable.Seq[FilterCriteria]]
+ val filtersByStreamId = new ConcurrentHashMap[String, Seq[FilterCriteria]]
}
def apply(
@@ -146,10 +145,10 @@ import org.slf4j.LoggerFactory
// The state must survive the actor lifecycle so keeping the state in an
Extension. Single writer per streamId.
private val storeExt = LocalConsumerFilterStore.StoreExt(context.system)
- def getState(): immutable.Seq[FilterCriteria] =
+ def getState(): Seq[FilterCriteria] =
storeExt.filtersByStreamId.computeIfAbsent(streamId, _ =>
Vector.empty[FilterCriteria])
- def setState(old: immutable.Seq[FilterCriteria], filterCriteria:
immutable.Seq[FilterCriteria]): Unit = {
+ def setState(old: Seq[FilterCriteria], filterCriteria: Seq[FilterCriteria]):
Unit = {
if (!storeExt.filtersByStreamId.replace(streamId, old, filterCriteria))
throw new ConcurrentModificationException(s"Unexpected concurrent update
of streamId [$streamId]")
context.log.debug2("Updated filter for streamId [{}] to [{}]", streamId,
filterCriteria)
@@ -207,7 +206,7 @@ import org.slf4j.LoggerFactory
// FIXME implement delta crdt
- def updated(filterCriteria: immutable.Seq[ConsumerFilter.FilterCriteria])(
+ def updated(filterCriteria: Seq[ConsumerFilter.FilterCriteria])(
implicit node: SelfUniqueAddress): State = {
var newExcludeTags = excludeTags
@@ -277,7 +276,7 @@ import org.slf4j.LoggerFactory
includeEntityOffsets = newIncludeEntityOffsets)
}
- lazy val toFilterCriteria: immutable.Seq[ConsumerFilter.FilterCriteria] = {
+ lazy val toFilterCriteria: Seq[ConsumerFilter.FilterCriteria] = {
Vector(
if (excludeTags.isEmpty) None
else Some(ConsumerFilter.ExcludeTags(excludeTags.elements)),
diff --git
a/grpc/src/main/scala/org/apache/pekko/projection/grpc/internal/ProtoAnySerialization.scala
b/grpc/src/main/scala/org/apache/pekko/projection/grpc/internal/ProtoAnySerialization.scala
index a9a0dc4f..9a95528f 100644
---
a/grpc/src/main/scala/org/apache/pekko/projection/grpc/internal/ProtoAnySerialization.scala
+++
b/grpc/src/main/scala/org/apache/pekko/projection/grpc/internal/ProtoAnySerialization.scala
@@ -14,7 +14,6 @@
package org.apache.pekko.projection.grpc.internal
import scala.collection.concurrent.TrieMap
-import scala.collection.immutable
import scala.jdk.CollectionConverters._
import scala.util.Try
@@ -106,7 +105,7 @@ import scalapb.options.Scalapb
*/
@InternalApi private[pekko] class ProtoAnySerialization(
system: ActorSystem[?],
- descriptors: immutable.Seq[Descriptors.FileDescriptor],
+ descriptors: Seq[Descriptors.FileDescriptor],
prefer: ProtoAnySerialization.Prefer) {
import ProtoAnySerialization._
diff --git
a/integration-examples/src/test/scala/docs/cassandra/CassandraProjectionDocExample.scala
b/integration-examples/src/test/scala/docs/cassandra/CassandraProjectionDocExample.scala
index c495fc78..383e342f 100644
---
a/integration-examples/src/test/scala/docs/cassandra/CassandraProjectionDocExample.scala
+++
b/integration-examples/src/test/scala/docs/cassandra/CassandraProjectionDocExample.scala
@@ -76,12 +76,10 @@ object CassandraProjectionDocExample {
// #handler
// #grouped-handler
- import scala.collection.immutable
-
- class GroupedShoppingCartHandler extends
Handler[immutable.Seq[EventEnvelope[ShoppingCart.Event]]] {
+ class GroupedShoppingCartHandler extends
Handler[Seq[EventEnvelope[ShoppingCart.Event]]] {
private val logger = LoggerFactory.getLogger(getClass)
- override def process(envelopes:
immutable.Seq[EventEnvelope[ShoppingCart.Event]]): Future[Done] = {
+ override def process(envelopes: Seq[EventEnvelope[ShoppingCart.Event]]):
Future[Done] = {
envelopes.map(_.event).foreach {
case ShoppingCart.CheckedOut(cartId, time) =>
logger.info2("Shopping cart {} was checked out at {}", cartId, time)
diff --git
a/jdbc-int-test/src/test/scala/org/apache/pekko/projection/jdbc/JdbcProjectionSpec.scala
b/jdbc-int-test/src/test/scala/org/apache/pekko/projection/jdbc/JdbcProjectionSpec.scala
index 45bd0914..3b2091c0 100644
---
a/jdbc-int-test/src/test/scala/org/apache/pekko/projection/jdbc/JdbcProjectionSpec.scala
+++
b/jdbc-int-test/src/test/scala/org/apache/pekko/projection/jdbc/JdbcProjectionSpec.scala
@@ -20,7 +20,6 @@ import java.util.concurrent.atomic.AtomicInteger
import java.util.concurrent.atomic.AtomicReference
import scala.annotation.tailrec
-import scala.collection.immutable
import scala.concurrent.Await
import scala.concurrent.ExecutionContext
import scala.concurrent.Future
@@ -631,7 +630,7 @@ class JdbcProjectionSpec
sourceProvider = sourceProvider(entityId),
jdbcSessionFactory,
handler = () =>
- JdbcHandler[PureJdbcSession, immutable.Seq[Envelope]] { (sess,
envelopes) =>
+ JdbcHandler[PureJdbcSession, Seq[Envelope]] { (sess, envelopes)
=>
handlerProbe.ref ! handlerCalled
sess.withConnection { conn =>
envelopes.foreach { envelope =>
@@ -658,8 +657,8 @@ class JdbcProjectionSpec
val result = new StringBuffer()
- def handler(): Handler[immutable.Seq[Envelope]] = new
Handler[immutable.Seq[Envelope]] {
- override def process(envelopes: immutable.Seq[Envelope]): Future[Done]
= {
+ def handler(): Handler[Seq[Envelope]] = new Handler[Seq[Envelope]] {
+ override def process(envelopes: Seq[Envelope]): Future[Done] = {
Future {
envelopes.foreach(env => result.append(env.message).append("|"))
}.map(_ => Done)
diff --git
a/jdbc/src/main/scala/org/apache/pekko/projection/jdbc/internal/Dialect.scala
b/jdbc/src/main/scala/org/apache/pekko/projection/jdbc/internal/Dialect.scala
index 0025814c..4a5f60a9 100644
---
a/jdbc/src/main/scala/org/apache/pekko/projection/jdbc/internal/Dialect.scala
+++
b/jdbc/src/main/scala/org/apache/pekko/projection/jdbc/internal/Dialect.scala
@@ -13,8 +13,6 @@
package org.apache.pekko.projection.jdbc.internal
-import scala.collection.immutable
-
import org.apache.pekko
import pekko.annotation.InternalApi
import pekko.util.Helpers.toRootLowerCase
@@ -40,7 +38,7 @@ private[projection] trait Dialect {
def tableName: String
def managementTableName: String
- def createTableStatements: immutable.Seq[String]
+ def createTableStatements: Seq[String]
def dropTableStatement: String
def readOffsetQuery: String
@@ -48,7 +46,7 @@ private[projection] trait Dialect {
def insertStatement(): String
def updateStatement(): String
- def createManagementTableStatements: immutable.Seq[String]
+ def createManagementTableStatements: Seq[String]
def dropManagementTableStatement: String
def readManagementStateQuery: String
def insertManagementStatement(): String
@@ -63,8 +61,8 @@ private[projection] trait Dialect {
@InternalApi
private[projection] object DialectDefaults {
- def createTableStatement(table: String): immutable.Seq[String] =
- immutable.Seq(
+ def createTableStatement(table: String): Seq[String] =
+ Seq(
s"""CREATE TABLE IF NOT EXISTS $table (
| "PROJECTION_NAME" VARCHAR(255) NOT NULL,
| "PROJECTION_KEY" VARCHAR(255) NOT NULL,
@@ -125,8 +123,8 @@ private[projection] object DialectDefaults {
val PROJECTION_KEY = 6
}
- def createManagementTableStatement(table: String): immutable.Seq[String] =
- immutable.Seq(s"""CREATE TABLE IF NOT EXISTS $table (
+ def createManagementTableStatement(table: String): Seq[String] =
+ Seq(s"""CREATE TABLE IF NOT EXISTS $table (
| "PROJECTION_NAME" VARCHAR(255) NOT NULL,
| "PROJECTION_KEY" VARCHAR(255) NOT NULL,
| "PAUSED" BOOLEAN NOT NULL,
@@ -196,7 +194,7 @@ private[projection] case class H2Dialect(
private val managementTable = transform(
schema.map(s =>
s""""$s"."$managementTableName"""").getOrElse(s""""$managementTableName""""))
- override val createTableStatements: immutable.Seq[String] =
+ override val createTableStatements: Seq[String] =
DialectDefaults.createTableStatement(table).map(s => transform(s))
override val dropTableStatement: String =
transform(DialectDefaults.dropTableStatement(table))
@@ -209,7 +207,7 @@ private[projection] case class H2Dialect(
override def updateStatement(): String =
transform(DialectDefaults.updateStatement(table))
- override val createManagementTableStatements: immutable.Seq[String] =
+ override val createManagementTableStatements: Seq[String] =
DialectDefaults.createManagementTableStatement(managementTable).map(s =>
transform(s))
override val dropManagementTableStatement: String = transform(
@@ -308,7 +306,7 @@ private[projection] case class MySQLDialect(schema:
Option[String], tableName: S
private val managementTable = schema.map(s =>
s"$s.$managementTableName").getOrElse(managementTableName)
override val createTableStatements =
- immutable.Seq(
+ Seq(
s"""CREATE TABLE IF NOT EXISTS $table (
| projection_name VARCHAR(255) NOT NULL,
| projection_key VARCHAR(255) NOT NULL,
@@ -337,7 +335,7 @@ private[projection] case class MySQLDialect(schema:
Option[String], tableName: S
Dialect.removeQuotes(DialectDefaults.updateStatement(table))
override val createManagementTableStatements =
- immutable.Seq(s"""CREATE TABLE IF NOT EXISTS $managementTable (
+ Seq(s"""CREATE TABLE IF NOT EXISTS $managementTable (
| projection_name VARCHAR(255) NOT NULL,
| projection_key VARCHAR(255) NOT NULL,
| paused BOOLEAN NOT NULL,
@@ -375,7 +373,7 @@ private[projection] case class MSSQLServerDialect(
private val managementTable = schema.map(s =>
s"""$s.$managementTableName""").getOrElse(s"""$managementTableName""")
override val createTableStatements =
- immutable.Seq(
+ Seq(
s"""IF NOT EXISTS (SELECT 1 FROM sys.objects WHERE object_id =
OBJECT_ID(N'$table') AND type in (N'U'))
|begin
| create table $table (
@@ -403,7 +401,7 @@ private[projection] case class MSSQLServerDialect(
override def updateStatement(): String =
DialectDefaults.updateStatement(table)
override val createManagementTableStatements =
- immutable.Seq(
+ Seq(
s"""IF NOT EXISTS (SELECT 1 FROM sys.objects WHERE object_id =
OBJECT_ID(N'$managementTable') AND type in (N'U'))
|begin
| create table $managementTable (
@@ -446,7 +444,7 @@ private[projection] case class OracleDialect(_schema:
Option[String], _tableName
schema.map(s =>
s""""$s"."$managementTableName"""").getOrElse(s""""$managementTableName"""")
override val createTableStatements =
- immutable.Seq(s"""
+ Seq(s"""
|BEGIN
|
| execute immediate 'create table $table ("PROJECTION_NAME"
VARCHAR2(255) NOT NULL,"PROJECTION_KEY" VARCHAR2(255) NOT NULL,"CURRENT_OFFSET"
VARCHAR2(255) NOT NULL,"MANIFEST" VARCHAR2(4) NOT NULL,"MERGEABLE" CHAR(1) NOT
NULL check ("MERGEABLE" in (0, 1)),"LAST_UPDATED" NUMBER(19) NOT NULL) ';
@@ -480,7 +478,7 @@ private[projection] case class OracleDialect(_schema:
Option[String], _tableName
override def updateStatement(): String =
DialectDefaults.updateStatement(table)
override val createManagementTableStatements =
- immutable.Seq(s"""
+ Seq(s"""
|BEGIN
|
| execute immediate 'create table $managementTable
("PROJECTION_NAME" VARCHAR2(255) NOT NULL,"PROJECTION_KEY" VARCHAR2(255) NOT
NULL,"PAUSED" CHAR(1) NOT NULL check ("PAUSED" in (0, 1)),"LAST_UPDATED"
NUMBER(19) NOT NULL) ';
diff --git
a/jdbc/src/main/scala/org/apache/pekko/projection/jdbc/internal/JdbcHandlerAdapter.scala
b/jdbc/src/main/scala/org/apache/pekko/projection/jdbc/internal/JdbcHandlerAdapter.scala
index e2b2fe13..80753c32 100644
---
a/jdbc/src/main/scala/org/apache/pekko/projection/jdbc/internal/JdbcHandlerAdapter.scala
+++
b/jdbc/src/main/scala/org/apache/pekko/projection/jdbc/internal/JdbcHandlerAdapter.scala
@@ -13,7 +13,6 @@
package org.apache.pekko.projection.jdbc.internal
-import scala.collection.immutable
import scala.jdk.CollectionConverters._
import org.apache.pekko
@@ -38,13 +37,13 @@ import pekko.projection.jdbc.scaladsl
}
/**
- * INTERNAL API: Adapter from `javadsl.Handler[java.util.List[Envelope]]` to
`scaladsl.Handler[immutable.Seq[Envelope]]`
+ * INTERNAL API: Adapter from `javadsl.Handler[java.util.List[Envelope]]` to
`scaladsl.Handler[Seq[Envelope]]`
*/
@InternalApi private[projection] class GroupedJdbcHandlerAdapter[Envelope, S
<: JdbcSession](
delegate: javadsl.JdbcHandler[java.util.List[Envelope], S])
- extends scaladsl.JdbcHandler[immutable.Seq[Envelope], S] {
+ extends scaladsl.JdbcHandler[Seq[Envelope], S] {
- override def process(session: S, envelopes: immutable.Seq[Envelope]): Unit =
{
+ override def process(session: S, envelopes: Seq[Envelope]): Unit = {
delegate.process(session, envelopes.asJava)
}
diff --git
a/jdbc/src/main/scala/org/apache/pekko/projection/jdbc/internal/JdbcProjectionImpl.scala
b/jdbc/src/main/scala/org/apache/pekko/projection/jdbc/internal/JdbcProjectionImpl.scala
index 1417c009..bf863e07 100644
---
a/jdbc/src/main/scala/org/apache/pekko/projection/jdbc/internal/JdbcProjectionImpl.scala
+++
b/jdbc/src/main/scala/org/apache/pekko/projection/jdbc/internal/JdbcProjectionImpl.scala
@@ -13,7 +13,6 @@
package org.apache.pekko.projection.jdbc.internal
-import scala.collection.immutable
import scala.concurrent.ExecutionContext
import scala.concurrent.Future
import scala.concurrent.duration.FiniteDuration
@@ -104,10 +103,10 @@ private[projection] object JdbcProjectionImpl {
projectionId: ProjectionId,
sourceProvider: SourceProvider[Offset, Envelope],
sessionFactory: () => S,
- handlerFactory: () => JdbcHandler[immutable.Seq[Envelope], S],
- offsetStore: JdbcOffsetStore[S]): () => Handler[immutable.Seq[Envelope]]
= { () =>
+ handlerFactory: () => JdbcHandler[Seq[Envelope], S],
+ offsetStore: JdbcOffsetStore[S]): () => Handler[Seq[Envelope]] = { () =>
new AdaptedJdbcHandler(handlerFactory(), offsetStore.executionContext) {
- override def process(envelopes: immutable.Seq[Envelope]): Future[Done] =
{
+ override def process(envelopes: Seq[Envelope]): Future[Done] = {
val offset = sourceProvider.extractOffset(envelopes.last)
JdbcSessionUtil
.withSession(sessionFactory) { sess =>
diff --git
a/jdbc/src/main/scala/org/apache/pekko/projection/jdbc/scaladsl/JdbcProjection.scala
b/jdbc/src/main/scala/org/apache/pekko/projection/jdbc/scaladsl/JdbcProjection.scala
index 587c8744..1d009c91 100644
---
a/jdbc/src/main/scala/org/apache/pekko/projection/jdbc/scaladsl/JdbcProjection.scala
+++
b/jdbc/src/main/scala/org/apache/pekko/projection/jdbc/scaladsl/JdbcProjection.scala
@@ -13,7 +13,6 @@
package org.apache.pekko.projection.jdbc.scaladsl
-import scala.collection.immutable
import scala.concurrent.Future
import scala.concurrent.duration.Duration
@@ -166,7 +165,7 @@ object JdbcProjection {
projectionId: ProjectionId,
sourceProvider: SourceProvider[Offset, Envelope],
sessionFactory: () => S,
- handler: () => JdbcHandler[immutable.Seq[Envelope], S])(
+ handler: () => JdbcHandler[Seq[Envelope], S])(
implicit system: ActorSystem[?]): GroupedProjection[Offset, Envelope] = {
val offsetStore = JdbcProjectionImpl.createOffsetStore(sessionFactory)
@@ -205,7 +204,7 @@ object JdbcProjection {
projectionId: ProjectionId,
sourceProvider: SourceProvider[Offset, Envelope],
sessionFactory: () => S,
- handler: () => Handler[immutable.Seq[Envelope]])(
+ handler: () => Handler[Seq[Envelope]])(
implicit system: ActorSystem[?]): GroupedProjection[Offset, Envelope] = {
val offsetStore = JdbcProjectionImpl.createOffsetStore(sessionFactory)
diff --git
a/kafka-test/src/test/scala/org/apache/pekko/projection/kafka/integration/KafkaToSlickIntegrationSpec.scala
b/kafka-test/src/test/scala/org/apache/pekko/projection/kafka/integration/KafkaToSlickIntegrationSpec.scala
index da5e046d..32963516 100644
---
a/kafka-test/src/test/scala/org/apache/pekko/projection/kafka/integration/KafkaToSlickIntegrationSpec.scala
+++
b/kafka-test/src/test/scala/org/apache/pekko/projection/kafka/integration/KafkaToSlickIntegrationSpec.scala
@@ -16,8 +16,6 @@ package org.apache.pekko.projection.kafka.integration
import java.lang.{ Long => JLong }
import java.util.concurrent.atomic.AtomicBoolean
-import scala.collection.immutable
-import scala.collection.immutable.Seq
import scala.concurrent.Await
import scala.concurrent.ExecutionContext
import scala.concurrent.Future
@@ -331,7 +329,7 @@ class KafkaToSlickIntegrationSpec extends
KafkaSpecBase(ConfigFactory.load().wit
}
}
- def produceEvents(topic: String, range: immutable.Seq[UserEvent], partition:
Int = 0): Future[Done] =
+ def produceEvents(topic: String, range: Seq[UserEvent], partition: Int = 0):
Future[Done] =
Source(range)
.map(e => new ProducerRecord(topic, partition, e.userId, e.eventType))
.runWith(Producer.plainSink(producerDefaults.withProducer(testProducer)))
diff --git
a/r2dbc-int-test/src/test/scala/org/apache/pekko/projection/r2dbc/R2dbcProjectionSpec.scala
b/r2dbc-int-test/src/test/scala/org/apache/pekko/projection/r2dbc/R2dbcProjectionSpec.scala
index a545510b..7c765ab9 100644
---
a/r2dbc-int-test/src/test/scala/org/apache/pekko/projection/r2dbc/R2dbcProjectionSpec.scala
+++
b/r2dbc-int-test/src/test/scala/org/apache/pekko/projection/r2dbc/R2dbcProjectionSpec.scala
@@ -18,7 +18,6 @@ import java.util.concurrent.atomic.AtomicInteger
import java.util.concurrent.atomic.AtomicReference
import scala.annotation.tailrec
-import scala.collection.immutable
import scala.concurrent.Await
import scala.concurrent.ExecutionContext
import scala.concurrent.Future
@@ -598,7 +597,7 @@ class R2dbcProjectionSpec
Some(settings),
sourceProvider = sourceProvider(entityId),
handler = () =>
- R2dbcHandler[immutable.Seq[Envelope]] { (session, envelopes) =>
+ R2dbcHandler[Seq[Envelope]] { (session, envelopes) =>
handlerProbe.ref ! handlerCalled
if (envelopes.isEmpty)
Future.successful(Done)
@@ -633,8 +632,8 @@ class R2dbcProjectionSpec
val result = new StringBuffer()
- def handler(): Handler[immutable.Seq[Envelope]] = new
Handler[immutable.Seq[Envelope]] {
- override def process(envelopes: immutable.Seq[Envelope]): Future[Done]
= {
+ def handler(): Handler[Seq[Envelope]] = new Handler[Seq[Envelope]] {
+ override def process(envelopes: Seq[Envelope]): Future[Done] = {
Future {
envelopes.foreach(env => result.append(env.message).append("|"))
}.map(_ => Done)
diff --git
a/r2dbc-int-test/src/test/scala/org/apache/pekko/projection/r2dbc/R2dbcTimestampOffsetProjectionSpec.scala
b/r2dbc-int-test/src/test/scala/org/apache/pekko/projection/r2dbc/R2dbcTimestampOffsetProjectionSpec.scala
index ed1f6cc3..cc5b022a 100644
---
a/r2dbc-int-test/src/test/scala/org/apache/pekko/projection/r2dbc/R2dbcTimestampOffsetProjectionSpec.scala
+++
b/r2dbc-int-test/src/test/scala/org/apache/pekko/projection/r2dbc/R2dbcTimestampOffsetProjectionSpec.scala
@@ -332,8 +332,8 @@ class R2dbcTimestampOffsetProjectionSpec
createEnvelope(pid1, 6, startTime.plusMillis(9), "e1-6"))
}
- def groupedHandler(probe: ActorRef[String]):
R2dbcHandler[immutable.Seq[EventEnvelope[String]]] = {
- R2dbcHandler[immutable.Seq[EventEnvelope[String]]] { (session, envelopes)
=>
+ def groupedHandler(probe: ActorRef[String]):
R2dbcHandler[Seq[EventEnvelope[String]]] = {
+ R2dbcHandler[Seq[EventEnvelope[String]]] { (session, envelopes) =>
probe ! "called"
if (envelopes.isEmpty)
Future.successful(Done)
@@ -745,9 +745,9 @@ class R2dbcTimestampOffsetProjectionSpec
val result = new StringBuffer()
- def handler(): Handler[immutable.Seq[EventEnvelope[String]]] =
- new Handler[immutable.Seq[EventEnvelope[String]]] {
- override def process(envelopes:
immutable.Seq[EventEnvelope[String]]): Future[Done] = {
+ def handler(): Handler[Seq[EventEnvelope[String]]] =
+ new Handler[Seq[EventEnvelope[String]]] {
+ override def process(envelopes: Seq[EventEnvelope[String]]):
Future[Done] = {
Future {
envelopes.foreach(env => result.append(env.event).append("|"))
}.map(_ => Done)
@@ -780,8 +780,8 @@ class R2dbcTimestampOffsetProjectionSpec
val result1 = new StringBuffer()
val result2 = new StringBuffer()
- def handler(): Handler[immutable.Seq[EventEnvelope[String]]] = new
Handler[immutable.Seq[EventEnvelope[String]]] {
- override def process(envelopes: immutable.Seq[EventEnvelope[String]]):
Future[Done] = {
+ def handler(): Handler[Seq[EventEnvelope[String]]] = new
Handler[Seq[EventEnvelope[String]]] {
+ override def process(envelopes: Seq[EventEnvelope[String]]):
Future[Done] = {
Future
.successful {
envelopes.foreach { envelope =>
@@ -820,8 +820,8 @@ class R2dbcTimestampOffsetProjectionSpec
val result1 = new StringBuffer()
val result2 = new StringBuffer()
- def handler(): Handler[immutable.Seq[EventEnvelope[String]]] = new
Handler[immutable.Seq[EventEnvelope[String]]] {
- override def process(envelopes: immutable.Seq[EventEnvelope[String]]):
Future[Done] = {
+ def handler(): Handler[Seq[EventEnvelope[String]]] = new
Handler[Seq[EventEnvelope[String]]] {
+ override def process(envelopes: Seq[EventEnvelope[String]]):
Future[Done] = {
Future
.successful {
envelopes.foreach { envelope =>
@@ -878,9 +878,9 @@ class R2dbcTimestampOffsetProjectionSpec
val result = new StringBuffer()
- def handler(): Handler[immutable.Seq[EventEnvelope[String]]] =
- new Handler[immutable.Seq[EventEnvelope[String]]] {
- override def process(envelopes:
immutable.Seq[EventEnvelope[String]]): Future[Done] = {
+ def handler(): Handler[Seq[EventEnvelope[String]]] =
+ new Handler[Seq[EventEnvelope[String]]] {
+ override def process(envelopes: Seq[EventEnvelope[String]]):
Future[Done] = {
Future {
envelopes.foreach(env => result.append(env.event).append("|"))
}.map(_ => Done)
diff --git
a/r2dbc/src/main/scala/org/apache/pekko/projection/r2dbc/internal/R2dbcHandlerAdapter.scala
b/r2dbc/src/main/scala/org/apache/pekko/projection/r2dbc/internal/R2dbcHandlerAdapter.scala
index 1028e2de..55a480b2 100644
---
a/r2dbc/src/main/scala/org/apache/pekko/projection/r2dbc/internal/R2dbcHandlerAdapter.scala
+++
b/r2dbc/src/main/scala/org/apache/pekko/projection/r2dbc/internal/R2dbcHandlerAdapter.scala
@@ -13,7 +13,6 @@
package org.apache.pekko.projection.r2dbc.internal
-import scala.collection.immutable
import scala.concurrent.Future
import scala.jdk.CollectionConverters._
import scala.jdk.FutureConverters._
@@ -45,13 +44,13 @@ import pekko.projection.r2dbc.scaladsl
/**
* INTERNAL API: Adapter from `javadsl.R2dbcHandler[java.util.List[Envelope]]`
to
- * `scaladsl.R2dbcHandler[immutable.Seq[Envelope]]`
+ * `scaladsl.R2dbcHandler[Seq[Envelope]]`
*/
@InternalApi private[projection] class R2dbcGroupedHandlerAdapter[Envelope](
delegate: javadsl.R2dbcHandler[java.util.List[Envelope]])
- extends scaladsl.R2dbcHandler[immutable.Seq[Envelope]] {
+ extends scaladsl.R2dbcHandler[Seq[Envelope]] {
- override def process(session: scaladsl.R2dbcSession, envelopes:
immutable.Seq[Envelope]): Future[Done] = {
+ override def process(session: scaladsl.R2dbcSession, envelopes:
Seq[Envelope]): Future[Done] = {
delegate.process(new R2dbcSession(session.connection)(session.ec,
session.system), envelopes.asJava).asScala
}
diff --git
a/r2dbc/src/main/scala/org/apache/pekko/projection/r2dbc/internal/R2dbcOffsetStore.scala
b/r2dbc/src/main/scala/org/apache/pekko/projection/r2dbc/internal/R2dbcOffsetStore.scala
index da9c426c..9ce55b0c 100644
---
a/r2dbc/src/main/scala/org/apache/pekko/projection/r2dbc/internal/R2dbcOffsetStore.scala
+++
b/r2dbc/src/main/scala/org/apache/pekko/projection/r2dbc/internal/R2dbcOffsetStore.scala
@@ -693,7 +693,7 @@ private[projection] class R2dbcOffsetStore(
case None => 0L
}
- def validateAll[Envelope](envelopes: immutable.Seq[Envelope]):
Future[immutable.Seq[(Envelope, Validation)]] = {
+ def validateAll[Envelope](envelopes: Seq[Envelope]): Future[Seq[(Envelope,
Validation)]] = {
import Validation._
envelopes
.foldLeft(Future.successful((getInflight(), Vector.empty[(Envelope,
Validation)]))) { (acc, envelope) =>
@@ -869,7 +869,7 @@ private[projection] class R2dbcOffsetStore(
}
}
- @tailrec final def addInflights[Envelope](envelopes:
immutable.Seq[Envelope]): Unit = {
+ @tailrec final def addInflights[Envelope](envelopes: Seq[Envelope]): Unit = {
val currentInflight = getInflight()
val entries = envelopes.iterator.map(createRecordWithOffset).collect {
case Some(r) =>
diff --git
a/r2dbc/src/main/scala/org/apache/pekko/projection/r2dbc/internal/R2dbcProjectionImpl.scala
b/r2dbc/src/main/scala/org/apache/pekko/projection/r2dbc/internal/R2dbcProjectionImpl.scala
index fefa9158..41f10257 100644
---
a/r2dbc/src/main/scala/org/apache/pekko/projection/r2dbc/internal/R2dbcProjectionImpl.scala
+++
b/r2dbc/src/main/scala/org/apache/pekko/projection/r2dbc/internal/R2dbcProjectionImpl.scala
@@ -16,7 +16,6 @@ package org.apache.pekko.projection.r2dbc.internal
import java.util.concurrent.atomic.AtomicLong
import scala.annotation.nowarn
-import scala.collection.immutable
import scala.concurrent.ExecutionContext
import scala.concurrent.Future
import scala.concurrent.duration.FiniteDuration
@@ -227,14 +226,14 @@ private[projection] object R2dbcProjectionImpl {
private[projection] def adaptedHandlerForGrouped[Offset, Envelope](
sourceProvider: SourceProvider[Offset, Envelope],
- handlerFactory: () => R2dbcHandler[immutable.Seq[Envelope]],
+ handlerFactory: () => R2dbcHandler[Seq[Envelope]],
offsetStore: R2dbcOffsetStore,
r2dbcExecutor: R2dbcExecutor)(
implicit
ec: ExecutionContext,
- system: ActorSystem[?]): () => Handler[immutable.Seq[Envelope]] = { () =>
+ system: ActorSystem[?]): () => Handler[Seq[Envelope]] = { () =>
new AdaptedR2dbcHandler(handlerFactory()) {
- override def process(envelopes: immutable.Seq[Envelope]): Future[Done] =
{
+ override def process(envelopes: Seq[Envelope]): Future[Done] = {
import R2dbcOffsetStore.Validation._
offsetStore.validateAll(envelopes).flatMap { isAcceptedEnvelopes =>
isAcceptedEnvelopes.foreach {
@@ -364,13 +363,13 @@ private[projection] object R2dbcProjectionImpl {
private[projection] def adaptedHandlerForGroupedAsync[Offset, Envelope](
sourceProvider: SourceProvider[Offset, Envelope],
- handlerFactory: () => Handler[immutable.Seq[Envelope]],
+ handlerFactory: () => Handler[Seq[Envelope]],
offsetStore: R2dbcOffsetStore)(
implicit
ec: ExecutionContext,
- system: ActorSystem[?]): () => Handler[immutable.Seq[Envelope]] = { () =>
+ system: ActorSystem[?]): () => Handler[Seq[Envelope]] = { () =>
new AdaptedHandler(handlerFactory()) {
- override def process(envelopes: immutable.Seq[Envelope]): Future[Done] =
{
+ override def process(envelopes: Seq[Envelope]): Future[Done] = {
import R2dbcOffsetStore.Validation._
offsetStore.validateAll(envelopes).flatMap { isAcceptedEnvelopes =>
isAcceptedEnvelopes.foreach {
@@ -673,7 +672,7 @@ private[projection] class R2dbcProjectionImpl[Offset,
Envelope](
override protected def saveOffsetsAndReport(
projectionId: ProjectionId,
- batch: immutable.Seq[ProjectionContextImpl[Offset, Envelope]]):
Future[Done] = {
+ batch: Seq[ProjectionContextImpl[Offset, Envelope]]): Future[Done] = {
import R2dbcProjectionImpl.FutureDone
val acceptedContexts =
diff --git
a/r2dbc/src/main/scala/org/apache/pekko/projection/r2dbc/scaladsl/R2dbcProjection.scala
b/r2dbc/src/main/scala/org/apache/pekko/projection/r2dbc/scaladsl/R2dbcProjection.scala
index 9417a890..9872a938 100644
---
a/r2dbc/src/main/scala/org/apache/pekko/projection/r2dbc/scaladsl/R2dbcProjection.scala
+++
b/r2dbc/src/main/scala/org/apache/pekko/projection/r2dbc/scaladsl/R2dbcProjection.scala
@@ -13,8 +13,6 @@
package org.apache.pekko.projection.r2dbc.scaladsl
-import scala.collection.immutable
-
import org.apache.pekko
import pekko.Done
import pekko.actor.typed.ActorSystem
@@ -226,7 +224,7 @@ object R2dbcProjection {
projectionId: ProjectionId,
settings: Option[R2dbcProjectionSettings],
sourceProvider: SourceProvider[Offset, Envelope],
- handler: () => R2dbcHandler[immutable.Seq[Envelope]])(implicit
+ handler: () => R2dbcHandler[Seq[Envelope]])(implicit
system: ActorSystem[?]): GroupedProjection[Offset, Envelope] =
groupedWithin(projectionId, ConfigFactory.empty(), settings,
sourceProvider, handler)
@@ -236,7 +234,7 @@ object R2dbcProjection {
config: Config,
settings: Option[R2dbcProjectionSettings],
sourceProvider: SourceProvider[Offset, Envelope],
- handler: () => R2dbcHandler[immutable.Seq[Envelope]])(implicit
+ handler: () => R2dbcHandler[Seq[Envelope]])(implicit
system: ActorSystem[?]): GroupedProjection[Offset, Envelope] = {
val r2dbcSettings = settings.getOrElse(R2dbcProjectionSettings(config,
system))
@@ -285,7 +283,7 @@ object R2dbcProjection {
projectionId: ProjectionId,
settings: Option[R2dbcProjectionSettings],
sourceProvider: SourceProvider[Offset, Envelope],
- handler: () => Handler[immutable.Seq[Envelope]])(implicit
+ handler: () => Handler[Seq[Envelope]])(implicit
system: ActorSystem[?]): GroupedProjection[Offset, Envelope] =
groupedWithinAsync(projectionId, ConfigFactory.empty(), settings,
sourceProvider, handler)
@@ -295,7 +293,7 @@ object R2dbcProjection {
config: Config,
settings: Option[R2dbcProjectionSettings],
sourceProvider: SourceProvider[Offset, Envelope],
- handler: () => Handler[immutable.Seq[Envelope]])(implicit
+ handler: () => Handler[Seq[Envelope]])(implicit
system: ActorSystem[?]): GroupedProjection[Offset, Envelope] = {
val r2dbcSettings = settings.getOrElse(R2dbcProjectionSettings(config,
system))
diff --git
a/r2dbc/src/test/scala/docs/home/projection/R2dbcProjectionDocExample.scala
b/r2dbc/src/test/scala/docs/home/projection/R2dbcProjectionDocExample.scala
index 2e0a1714..24e2ba60 100644
--- a/r2dbc/src/test/scala/docs/home/projection/R2dbcProjectionDocExample.scala
+++ b/r2dbc/src/test/scala/docs/home/projection/R2dbcProjectionDocExample.scala
@@ -80,15 +80,13 @@ object R2dbcProjectionDocExample {
// #handler
// #grouped-handler
- import scala.collection.immutable
-
class GroupedShoppingCartHandler()(implicit ec: ExecutionContext)
- extends R2dbcHandler[immutable.Seq[EventEnvelope[ShoppingCart.Event]]] {
+ extends R2dbcHandler[Seq[EventEnvelope[ShoppingCart.Event]]] {
private val logger = LoggerFactory.getLogger(getClass)
override def process(
session: R2dbcSession,
- envelopes: immutable.Seq[EventEnvelope[ShoppingCart.Event]]):
Future[Done] = {
+ envelopes: Seq[EventEnvelope[ShoppingCart.Event]]): Future[Done] = {
// save all events in DB
val stmts = envelopes
diff --git
a/slick/src/main/scala/org/apache/pekko/projection/slick/SlickProjection.scala
b/slick/src/main/scala/org/apache/pekko/projection/slick/SlickProjection.scala
index 758c5b06..b716b7bd 100644
---
a/slick/src/main/scala/org/apache/pekko/projection/slick/SlickProjection.scala
+++
b/slick/src/main/scala/org/apache/pekko/projection/slick/SlickProjection.scala
@@ -13,7 +13,6 @@
package org.apache.pekko.projection.slick
-import scala.collection.immutable
import scala.concurrent.ExecutionContext
import scala.concurrent.Future
import scala.concurrent.duration.Duration
@@ -227,20 +226,20 @@ object SlickProjection {
projectionId: ProjectionId,
sourceProvider: SourceProvider[Offset, Envelope],
databaseConfig: DatabaseConfig[P],
- handler: () => SlickHandler[immutable.Seq[Envelope]])(
+ handler: () => SlickHandler[Seq[Envelope]])(
implicit system: ActorSystem[?]): GroupedProjection[Offset, Envelope] = {
val offsetStore = createOffsetStore(databaseConfig)
- val adaptedSlickHandler: () => Handler[immutable.Seq[Envelope]] = () =>
- new Handler[immutable.Seq[Envelope]] {
+ val adaptedSlickHandler: () => Handler[Seq[Envelope]] = () =>
+ new Handler[Seq[Envelope]] {
import databaseConfig.profile.api._
private implicit val ec: ExecutionContext = system.executionContext
private val logger = Logging(system.classicSystem,
classOf[SlickProjectionImpl[?, ?, ?]])
private val delegate = handler()
- override def process(envelopes: immutable.Seq[Envelope]): Future[Done]
= {
+ override def process(envelopes: Seq[Envelope]): Future[Done] = {
val lastOffset = sourceProvider.extractOffset(envelopes.last)
val processedDBIO = offsetStore
@@ -308,7 +307,7 @@ object SlickProjection {
projectionId: ProjectionId,
sourceProvider: SourceProvider[Offset, Envelope],
databaseConfig: DatabaseConfig[P],
- handler: () => Handler[immutable.Seq[Envelope]])(
+ handler: () => Handler[Seq[Envelope]])(
implicit system: ActorSystem[?]): GroupedProjection[Offset, Envelope] = {
val offsetStore = createOffsetStore(databaseConfig)
diff --git
a/slick/src/test/scala/org/apache/pekko/projection/slick/SlickProjectionSpec.scala
b/slick/src/test/scala/org/apache/pekko/projection/slick/SlickProjectionSpec.scala
index 7fd1a426..d85d1aee 100644
---
a/slick/src/test/scala/org/apache/pekko/projection/slick/SlickProjectionSpec.scala
+++
b/slick/src/test/scala/org/apache/pekko/projection/slick/SlickProjectionSpec.scala
@@ -19,7 +19,6 @@ import java.util.concurrent.atomic.AtomicInteger
import java.util.concurrent.atomic.AtomicReference
import scala.annotation.tailrec
-import scala.collection.immutable
import scala.concurrent.Await
import scala.concurrent.ExecutionContext
import scala.concurrent.Future
@@ -702,7 +701,7 @@ class SlickProjectionSpec
databaseConfig = dbConfig,
// build event handler from simple lambda
handler = () =>
- SlickHandler[immutable.Seq[Envelope]] { envelopes =>
+ SlickHandler[Seq[Envelope]] { envelopes =>
handlerProbe.ref ! handlerCalled
val dbios = envelopes.map(env =>
repository.concatToText(env.id, env.message))
DBIOAction.sequence(dbios).map(_ => Done)
@@ -734,8 +733,8 @@ class SlickProjectionSpec
val result = new StringBuffer()
- def handler(): Handler[immutable.Seq[Envelope]] = new
Handler[immutable.Seq[Envelope]] {
- override def process(envelopes: immutable.Seq[Envelope]): Future[Done]
= {
+ def handler(): Handler[Seq[Envelope]] = new Handler[Seq[Envelope]] {
+ override def process(envelopes: Seq[Envelope]): Future[Done] = {
Future {
envelopes.foreach(env => result.append(env.message).append("|"))
}.map(_ => Done)
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]