This is an automated email from the ASF dual-hosted git repository.
He-Pin pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/pekko-connectors.git
The following commit(s) were added to refs/heads/main by this push:
new 5dd352b05 fix: fix compilation warnings across all Scala versions
(#1750)
5dd352b05 is described below
commit 5dd352b0538c3dab37214465b2b0ab3980d2ff0b
Author: He-Pin(kerr) <[email protected]>
AuthorDate: Mon Jul 6 16:23:55 2026 +0800
fix: fix compilation warnings across all Scala versions (#1750)
* fix: fix compilation warnings across all Scala versions
Motivation:
Compilation produced warnings across Scala 2.13 and Scala 3 including:
- Unused imports in google-cloud-pub-sub, google-fcm, geode, s3
- Pattern match exhaustivity in mqtt-streaming (Scala 3 E029)
- Anonymous class duplication at inline site in geode (Scala 3 E197)
- Deprecated API usage in aws-spi-pekko-http, kinesis, solr
- Type inferred to Object/Any in amqp, influxdb, jms, jakartams, s3
- Name shadowing in jms and jakartams test specs
- Structural type inference in kinesis Valve
- Suspicious unqualified getClass call in google-cloud-storage
- Unchecked calls in AmqpFlowTest Java tests
- Pattern type more specialized than RHS in google-cloud-pub-sub-grpc
Modification:
- Remove unused imports across multiple modules
- Add @nowarn annotations for intentional deprecated API usage and
pattern match exhaustivity in MQTT session event flows
- Fix pattern match in MqttFrameStage to use Left(_) instead of
unchecked type pattern
- Add @nowarn for anonymous class duplication in LabelledGenericGeneric
- Replace deprecated Stream with LazyList in kinesis tests
- Add explicit type annotations (Map[String, Any], [Any]) to fix
type inference warnings
- Rename shadowed pattern variables in JMS test specs
- Add explicit ValveSwitch type annotation in kinesis Valve
- Use this.getClass in GCStorageWiremockBase constructor
- Add @SuppressWarnings("unchecked") to Java test methods
- Update @nowarn filters in google-cloud-pub-sub-grpc ExampleApp
- Change .recover to .recoverWith in amqp test for correct Future chaining
Result:
Clean compilation with zero warnings across all Scala versions.
Tests:
sbt +test:compile - all versions pass with no warnings
References:
None - code quality improvement
* fix: restore private modifier on GCStorageWiremockBase.config
Address review feedback: the private modifier removal was unnecessary
since companion class/object can access each other's private members.
---
.../stream/connectors/amqp/javadsl/AmqpFlowTest.java | 3 +++
.../AmqpGraphStageLogicConnectionShutdownSpec.scala | 2 +-
.../pekko/stream/connectors/awsspi/s3/TestS3.scala | 2 ++
.../geode/impl/pdx/LabelledGenericGeneric.scala | 2 ++
geode/src/test/scala/docs/scaladsl/GeodeBaseSpec.scala | 1 -
.../scala/docs/scaladsl/GeodeContinuousSourceSpec.scala | 1 -
.../test/scala/docs/scaladsl/GeodeFiniteSourceSpec.scala | 1 -
geode/src/test/scala/docs/scaladsl/GeodeFlowSpec.scala | 1 -
.../src/test/scala/docs/scaladsl/ExampleApp.scala | 16 ++++++++--------
.../googlecloud/pubsub/scaladsl/GooglePubSub.scala | 2 --
.../storage/scaladsl/GCStorageWiremockBase.scala | 2 +-
.../connectors/google/firebase/fcm/FcmSettings.scala | 4 ----
.../test/scala/docs/scaladsl/InfluxDbSourceSpec.scala | 6 +++---
.../src/test/scala/docs/scaladsl/JmsConnectorsSpec.scala | 2 +-
.../jakartams/impl/JmsMessageProducerSpec.scala | 4 ++--
jms/src/test/scala/docs/scaladsl/JmsConnectorsSpec.scala | 2 +-
.../connectors/jms/impl/JmsMessageProducerSpec.scala | 4 ++--
.../stream/connectors/kinesis/KinesisFlowSpec.scala | 6 ++----
.../apache/pekko/stream/connectors/kinesis/Valve.scala | 4 +++-
.../connectors/mqtt/streaming/impl/MqttFrameStage.scala | 2 +-
.../connectors/mqtt/streaming/scaladsl/MqttSession.scala | 3 +++
.../pekko/stream/connectors/s3/scaladsl/Generators.scala | 1 -
.../pekko/stream/connectors/s3/scaladsl/S3ExtSpec.scala | 2 +-
.../connectors/s3/scaladsl/S3IntegrationSpec.scala | 6 +++---
solr/src/test/scala/docs/scaladsl/SolrSpec.scala | 2 ++
25 files changed, 41 insertions(+), 40 deletions(-)
diff --git
a/amqp/src/test/java/org/apache/pekko/stream/connectors/amqp/javadsl/AmqpFlowTest.java
b/amqp/src/test/java/org/apache/pekko/stream/connectors/amqp/javadsl/AmqpFlowTest.java
index 4629895dd..08045fed5 100644
---
a/amqp/src/test/java/org/apache/pekko/stream/connectors/amqp/javadsl/AmqpFlowTest.java
+++
b/amqp/src/test/java/org/apache/pekko/stream/connectors/amqp/javadsl/AmqpFlowTest.java
@@ -82,6 +82,7 @@ public class AmqpFlowTest {
AmqpFlow.createWithConfirmUnordered(settings(reuseByteArray)));
}
+ @SuppressWarnings("unchecked")
private void shouldEmitConfirmationForPublishedMessages(
final Flow<WriteMessage, WriteResult, CompletionStage<Done>> flow) {
@@ -113,6 +114,7 @@ public class AmqpFlowTest {
shouldPropagateContext(AmqpFlowWithContext.createWithConfirm(settings(reuseByteArray)));
}
+ @SuppressWarnings("unchecked")
private void shouldPropagateContext(
FlowWithContext<WriteMessage, String, WriteResult, String,
CompletionStage<Done>>
flowWithContext) {
@@ -137,6 +139,7 @@ public class AmqpFlowTest {
@ParameterizedTest
@ValueSource(booleans = {false, true})
+ @SuppressWarnings("unchecked")
public void shouldPropagatePassThrough(boolean reuseByteArray) {
Flow<Pair<WriteMessage, String>, Pair<WriteResult, String>,
CompletionStage<Done>> flow =
AmqpFlow.createWithConfirmAndPassThroughUnordered(settings(reuseByteArray));
diff --git
a/amqp/src/test/scala/org/apache/pekko/stream/connectors/amqp/scaladsl/AmqpGraphStageLogicConnectionShutdownSpec.scala
b/amqp/src/test/scala/org/apache/pekko/stream/connectors/amqp/scaladsl/AmqpGraphStageLogicConnectionShutdownSpec.scala
index ae249678f..08d6e19e6 100644
---
a/amqp/src/test/scala/org/apache/pekko/stream/connectors/amqp/scaladsl/AmqpGraphStageLogicConnectionShutdownSpec.scala
+++
b/amqp/src/test/scala/org/apache/pekko/stream/connectors/amqp/scaladsl/AmqpGraphStageLogicConnectionShutdownSpec.scala
@@ -97,7 +97,7 @@ class AmqpGraphStageLogicConnectionShutdownSpec
Future
.traverse(input)(in => Source.single(ByteString(in)).runWith(amqpSink))
- .recover {
+ .recoverWith {
case NonFatal(e) => system.terminate().flatMap(_ => Future.failed(e))
}
.flatMap(_ => system.terminate())
diff --git
a/aws-spi-pekko-http/src/test/scala/org/apache/pekko/stream/connectors/awsspi/s3/TestS3.scala
b/aws-spi-pekko-http/src/test/scala/org/apache/pekko/stream/connectors/awsspi/s3/TestS3.scala
index 28241f4ba..9f6b3f628 100644
---
a/aws-spi-pekko-http/src/test/scala/org/apache/pekko/stream/connectors/awsspi/s3/TestS3.scala
+++
b/aws-spi-pekko-http/src/test/scala/org/apache/pekko/stream/connectors/awsspi/s3/TestS3.scala
@@ -27,6 +27,7 @@ import software.amazon.awssdk.core.async.{ AsyncRequestBody,
AsyncResponseTransf
import software.amazon.awssdk.services.s3.{ S3AsyncClient, S3Configuration }
import software.amazon.awssdk.services.s3.model._
+import scala.annotation.nowarn
import scala.concurrent.duration._
import scala.util.Random
@@ -154,6 +155,7 @@ class TestS3 extends BaseAwsClientTest[S3AsyncClient] {
bucketName
}
+ @nowarn("msg=deprecated")
private def withClient(testCode: S3AsyncClient => Any): Any = {
val pekkoClient = new
PekkoHttpAsyncHttpService().createAsyncHttpClientFactory().build()
diff --git
a/geode/src/main/scala-3/org/apache/pekko/stream/connectors/geode/impl/pdx/LabelledGenericGeneric.scala
b/geode/src/main/scala-3/org/apache/pekko/stream/connectors/geode/impl/pdx/LabelledGenericGeneric.scala
index b98fd7a17..32e6803ad 100644
---
a/geode/src/main/scala-3/org/apache/pekko/stream/connectors/geode/impl/pdx/LabelledGenericGeneric.scala
+++
b/geode/src/main/scala-3/org/apache/pekko/stream/connectors/geode/impl/pdx/LabelledGenericGeneric.scala
@@ -19,6 +19,7 @@ package org.apache.pekko.stream.connectors.geode.impl.pdx
import org.apache.pekko.annotation.InternalApi
+import scala.annotation.nowarn
import scala.deriving.Mirror
@InternalApi
@@ -53,6 +54,7 @@ private[pekko] object LabelledGeneric {
inline def apply[A](using l: LabelledGeneric[A]): LabelledGeneric.Aux[A,
l.Repr] = l
+ @nowarn("msg=New anonymous class definition will be duplicated")
transparent inline given productInst[A <: Product](
using m: Mirror.ProductOf[A])
: LabelledGeneric.Aux[A, ZipWith[m.MirroredElemLabels,
m.MirroredElemTypes, FieldType]] =
diff --git a/geode/src/test/scala/docs/scaladsl/GeodeBaseSpec.scala
b/geode/src/test/scala/docs/scaladsl/GeodeBaseSpec.scala
index a4a3b6892..ff032eec0 100644
--- a/geode/src/test/scala/docs/scaladsl/GeodeBaseSpec.scala
+++ b/geode/src/test/scala/docs/scaladsl/GeodeBaseSpec.scala
@@ -24,7 +24,6 @@ import org.scalatest.BeforeAndAfterAll
import scala.concurrent.Await
import scala.concurrent.duration.DurationInt
-import scala.language.postfixOps
import org.scalatest.matchers.should.Matchers
import org.scalatest.wordspec.AnyWordSpec
diff --git a/geode/src/test/scala/docs/scaladsl/GeodeContinuousSourceSpec.scala
b/geode/src/test/scala/docs/scaladsl/GeodeContinuousSourceSpec.scala
index 3d46a91c1..1765a970c 100644
--- a/geode/src/test/scala/docs/scaladsl/GeodeContinuousSourceSpec.scala
+++ b/geode/src/test/scala/docs/scaladsl/GeodeContinuousSourceSpec.scala
@@ -21,7 +21,6 @@ import org.slf4j.LoggerFactory
import scala.concurrent.Await
import scala.concurrent.duration.DurationInt
-import scala.language.postfixOps
class GeodeContinuousSourceSpec extends GeodeBaseSpec {
diff --git a/geode/src/test/scala/docs/scaladsl/GeodeFiniteSourceSpec.scala
b/geode/src/test/scala/docs/scaladsl/GeodeFiniteSourceSpec.scala
index ea80ae2f8..5d042093d 100644
--- a/geode/src/test/scala/docs/scaladsl/GeodeFiniteSourceSpec.scala
+++ b/geode/src/test/scala/docs/scaladsl/GeodeFiniteSourceSpec.scala
@@ -20,7 +20,6 @@ import org.slf4j.LoggerFactory
import scala.concurrent.Await
import scala.concurrent.duration.DurationInt
-import scala.language.postfixOps
class GeodeFiniteSourceSpec extends GeodeBaseSpec {
diff --git a/geode/src/test/scala/docs/scaladsl/GeodeFlowSpec.scala
b/geode/src/test/scala/docs/scaladsl/GeodeFlowSpec.scala
index fd255386e..0f8a7a542 100644
--- a/geode/src/test/scala/docs/scaladsl/GeodeFlowSpec.scala
+++ b/geode/src/test/scala/docs/scaladsl/GeodeFlowSpec.scala
@@ -21,7 +21,6 @@ import pekko.stream.scaladsl.{ Flow, Sink }
import scala.concurrent.Await
import scala.concurrent.duration.DurationInt
-import scala.language.postfixOps
class GeodeFlowSpec extends GeodeBaseSpec {
diff --git
a/google-cloud-pub-sub-grpc/src/test/scala/docs/scaladsl/ExampleApp.scala
b/google-cloud-pub-sub-grpc/src/test/scala/docs/scaladsl/ExampleApp.scala
index 32417c6dc..ee84eab43 100644
--- a/google-cloud-pub-sub-grpc/src/test/scala/docs/scaladsl/ExampleApp.scala
+++ b/google-cloud-pub-sub-grpc/src/test/scala/docs/scaladsl/ExampleApp.scala
@@ -76,7 +76,7 @@ object ExampleApp {
}
private def publishSingle(args: List[String])(implicit system: ActorSystem)
= {
- val projectId :: topic :: Nil = args: @nowarn("msg=match may not be
exhaustive")
+ val projectId :: topic :: Nil = args: @nowarn("msg=match may not be
exhaustive|more specialized")
Source
.single(publish(projectId, topic)("Hello!"))
@@ -85,7 +85,7 @@ object ExampleApp {
}
private def publishStream(args: List[String])(implicit system: ActorSystem)
= {
- val projectId :: topic :: Nil = args: @nowarn("msg=match may not be
exhaustive")
+ val projectId :: topic :: Nil = args: @nowarn("msg=match may not be
exhaustive|more specialized")
Source
.tick(0.seconds, 1.second, ())
@@ -102,7 +102,7 @@ object ExampleApp {
}
private def subscribeStream(args: List[String])(implicit system:
ActorSystem) = {
- val projectId :: sub :: Nil = args: @nowarn("msg=match may not be
exhaustive")
+ val projectId :: sub :: Nil = args: @nowarn("msg=match may not be
exhaustive|more specialized")
GooglePubSub
.subscribe(subscribe(projectId, sub), 1.second)
@@ -111,7 +111,7 @@ object ExampleApp {
}
private def subscribeAutoExtend(args: List[String])(implicit system:
ActorSystem) = {
- val projectId :: sub :: Nil = args: @nowarn("msg=match may not be
exhaustive")
+ val projectId :: sub :: Nil = args: @nowarn("msg=match may not be
exhaustive|more specialized")
val subscriptionFqrs = subFqrs(projectId, sub)
val restartSettings = RestartSettings(
@@ -138,7 +138,7 @@ object ExampleApp {
* Usage: subscribe-flow-control <projectId> <subscription>
*/
private def subscribeFlowControl(args: List[String])(implicit system:
ActorSystem) = {
- val projectId :: sub :: Nil = args: @nowarn("msg=match may not be
exhaustive")
+ val projectId :: sub :: Nil = args: @nowarn("msg=match may not be
exhaustive|more specialized")
val subscriptionFqrs = subFqrs(projectId, sub)
val restartSettings = RestartSettings(
@@ -170,7 +170,7 @@ object ExampleApp {
private def subscribeNackRetry(args: List[String])(implicit system:
ActorSystem) = {
import system.dispatcher
- val projectId :: sub :: Nil = args: @nowarn("msg=match may not be
exhaustive")
+ val projectId :: sub :: Nil = args: @nowarn("msg=match may not be
exhaustive|more specialized")
val subscriptionFqrs = subFqrs(projectId, sub)
val restartSettings = RestartSettings(
@@ -206,7 +206,7 @@ object ExampleApp {
* Usage: subscribe-dynamic-deadline <projectId> <subscription>
*/
private def subscribeDynamicDeadline(args: List[String])(implicit system:
ActorSystem) = {
- val projectId :: sub :: Nil = args: @nowarn("msg=match may not be
exhaustive")
+ val projectId :: sub :: Nil = args: @nowarn("msg=match may not be
exhaustive|more specialized")
val subscriptionFqrs = subFqrs(projectId, sub)
val restartSettings = RestartSettings(
@@ -246,7 +246,7 @@ object ExampleApp {
private def subscribeAdaptive(args: List[String])(implicit system:
ActorSystem) = {
import
org.apache.pekko.stream.connectors.googlecloud.pubsub.grpc.AckDeadlineDistribution
- val projectId :: sub :: Nil = args: @nowarn("msg=match may not be
exhaustive")
+ val projectId :: sub :: Nil = args: @nowarn("msg=match may not be
exhaustive|more specialized")
val subscriptionFqrs = subFqrs(projectId, sub)
val restartSettings = RestartSettings(
diff --git
a/google-cloud-pub-sub/src/main/scala/org/apache/pekko/stream/connectors/googlecloud/pubsub/scaladsl/GooglePubSub.scala
b/google-cloud-pub-sub/src/main/scala/org/apache/pekko/stream/connectors/googlecloud/pubsub/scaladsl/GooglePubSub.scala
index 5ae4b99be..b0d7a7ad2 100644
---
a/google-cloud-pub-sub/src/main/scala/org/apache/pekko/stream/connectors/googlecloud/pubsub/scaladsl/GooglePubSub.scala
+++
b/google-cloud-pub-sub/src/main/scala/org/apache/pekko/stream/connectors/googlecloud/pubsub/scaladsl/GooglePubSub.scala
@@ -15,8 +15,6 @@ package
org.apache.pekko.stream.connectors.googlecloud.pubsub.scaladsl
import org.apache.pekko
import pekko.actor.Cancellable
-import pekko.stream.Attributes
-import pekko.stream.connectors.google.GoogleAttributes
import pekko.stream.connectors.googlecloud.pubsub._
import pekko.stream.connectors.googlecloud.pubsub.impl._
import pekko.stream.scaladsl.{ Flow, FlowWithContext, Keep, Sink, Source }
diff --git
a/google-cloud-storage/src/test/scala/org/apache/pekko/stream/connectors/googlecloud/storage/scaladsl/GCStorageWiremockBase.scala
b/google-cloud-storage/src/test/scala/org/apache/pekko/stream/connectors/googlecloud/storage/scaladsl/GCStorageWiremockBase.scala
index a454a8fae..c6de317ac 100644
---
a/google-cloud-storage/src/test/scala/org/apache/pekko/stream/connectors/googlecloud/storage/scaladsl/GCStorageWiremockBase.scala
+++
b/google-cloud-storage/src/test/scala/org/apache/pekko/stream/connectors/googlecloud/storage/scaladsl/GCStorageWiremockBase.scala
@@ -31,7 +31,7 @@ import scala.util.Random
abstract class GCStorageWiremockBase(_system: ActorSystem, _wireMockServer:
Hoverfly) extends TestKit(_system) {
def this(mock: Hoverfly) =
- this(ActorSystem(getCallerName(getClass),
+ this(ActorSystem(getCallerName(classOf[GCStorageWiremockBase]),
config(mock.getHoverflyConfig.getProxyPort).withFallback(ConfigFactory.load())),
mock)
diff --git
a/google-fcm/src/main/scala/org/apache/pekko/stream/connectors/google/firebase/fcm/FcmSettings.scala
b/google-fcm/src/main/scala/org/apache/pekko/stream/connectors/google/firebase/fcm/FcmSettings.scala
index 5ca35505e..e95363510 100644
---
a/google-fcm/src/main/scala/org/apache/pekko/stream/connectors/google/firebase/fcm/FcmSettings.scala
+++
b/google-fcm/src/main/scala/org/apache/pekko/stream/connectors/google/firebase/fcm/FcmSettings.scala
@@ -13,10 +13,6 @@
package org.apache.pekko.stream.connectors.google.firebase.fcm
-import org.apache.pekko
-
-import java.util.Objects
-
final class FcmSettings private (
val isTest: Boolean,
val maxConcurrentConnections: Int) {
diff --git a/influxdb/src/test/scala/docs/scaladsl/InfluxDbSourceSpec.scala
b/influxdb/src/test/scala/docs/scaladsl/InfluxDbSourceSpec.scala
index 45ebea21b..b09fe06ad 100644
--- a/influxdb/src/test/scala/docs/scaladsl/InfluxDbSourceSpec.scala
+++ b/influxdb/src/test/scala/docs/scaladsl/InfluxDbSourceSpec.scala
@@ -71,7 +71,7 @@ class InfluxDbSourceSpec
val query = new Query("SELECT man() FROM invalid", DatabaseName)
val result = InfluxDbSource(influxDB, query) // .runWith(Sink.seq)
- .recover {
+ .recover[Any] {
case e: InfluxDBException => e.getMessage
}
.runWith(Sink.seq)
@@ -98,7 +98,7 @@ class InfluxDbSourceSpec
val result = InfluxDbSource
.typed(classOf[InfluxDbSourceCpu], InfluxDbReadSettings.Default,
influxDB, query) // .runWith(Sink.seq)
- .recover {
+ .recover[Any] {
case e: InfluxDBException => e.getMessage
}
.runWith(Sink.seq)
@@ -112,7 +112,7 @@ class InfluxDbSourceSpec
val result = InfluxDbSource
.typed(classOf[InfluxDbSourceCpu], InfluxDbReadSettings.Default,
influxDB, query) // .runWith(Sink.seq)
- .recover {
+ .recover[Any] {
case e: InfluxDBException => e.getMessage
}
.runWith(Sink.seq)
diff --git a/jakartams/src/test/scala/docs/scaladsl/JmsConnectorsSpec.scala
b/jakartams/src/test/scala/docs/scaladsl/JmsConnectorsSpec.scala
index c4cbc94cd..bf9c647a6 100644
--- a/jakartams/src/test/scala/docs/scaladsl/JmsConnectorsSpec.scala
+++ b/jakartams/src/test/scala/docs/scaladsl/JmsConnectorsSpec.scala
@@ -1185,7 +1185,7 @@ class JmsConnectorsSpec extends JmsSpec {
val respondStreamControl: JmsConsumerControl =
JmsConsumer(JmsConsumerSettings(system,
connectionFactory).withQueue("test"))
.collect {
- case message: TextMessage => JmsTextMessage(message)
+ case textMsg: TextMessage => JmsTextMessage(textMsg)
}
.map { textMessage =>
textMessage.headers.foldLeft(JmsTextMessage(textMessage.body.reverse)) {
diff --git
a/jakartams/src/test/scala/org/apache/pekko/stream/connectors/jakartams/impl/JmsMessageProducerSpec.scala
b/jakartams/src/test/scala/org/apache/pekko/stream/connectors/jakartams/impl/JmsMessageProducerSpec.scala
index aed282a50..2e72947b4 100644
---
a/jakartams/src/test/scala/org/apache/pekko/stream/connectors/jakartams/impl/JmsMessageProducerSpec.scala
+++
b/jakartams/src/test/scala/org/apache/pekko/stream/connectors/jakartams/impl/JmsMessageProducerSpec.scala
@@ -68,7 +68,7 @@ class JmsMessageProducerSpec extends JmsSpec with
MockitoSugar {
}
"succeed if properties are set as map" in new Setup {
- val props = Map(
+ val props: Map[String, Any] = Map(
"string" -> "string",
"int" -> 1,
"boolean" -> true,
@@ -119,7 +119,7 @@ class JmsMessageProducerSpec extends JmsSpec with
MockitoSugar {
val jmsProducer = JmsMessageProducer(jmsSession, settings, 0)
jmsProducer.createMessage(
JmsMapMessage(
- Map(
+ Map[String, Any](
"string" -> "string",
"int" -> 1,
"boolean" -> true,
diff --git a/jms/src/test/scala/docs/scaladsl/JmsConnectorsSpec.scala
b/jms/src/test/scala/docs/scaladsl/JmsConnectorsSpec.scala
index 42c35f777..428308c51 100644
--- a/jms/src/test/scala/docs/scaladsl/JmsConnectorsSpec.scala
+++ b/jms/src/test/scala/docs/scaladsl/JmsConnectorsSpec.scala
@@ -1176,7 +1176,7 @@ class JmsConnectorsSpec extends JmsSpec {
val respondStreamControl: JmsConsumerControl =
JmsConsumer(JmsConsumerSettings(system,
connectionFactory).withQueue("test"))
.collect {
- case message: TextMessage => JmsTextMessage(message)
+ case textMsg: TextMessage => JmsTextMessage(textMsg)
}
.map { textMessage =>
textMessage.headers.foldLeft(JmsTextMessage(textMessage.body.reverse)) {
diff --git
a/jms/src/test/scala/org/apache/pekko/stream/connectors/jms/impl/JmsMessageProducerSpec.scala
b/jms/src/test/scala/org/apache/pekko/stream/connectors/jms/impl/JmsMessageProducerSpec.scala
index 5bf389d38..8182bcb5f 100644
---
a/jms/src/test/scala/org/apache/pekko/stream/connectors/jms/impl/JmsMessageProducerSpec.scala
+++
b/jms/src/test/scala/org/apache/pekko/stream/connectors/jms/impl/JmsMessageProducerSpec.scala
@@ -67,7 +67,7 @@ class JmsMessageProducerSpec extends JmsSpec with
MockitoSugar {
}
"succeed if properties are set as map" in new Setup {
- val props = Map(
+ val props: Map[String, Any] = Map(
"string" -> "string",
"int" -> 1,
"boolean" -> true,
@@ -118,7 +118,7 @@ class JmsMessageProducerSpec extends JmsSpec with
MockitoSugar {
val jmsProducer = JmsMessageProducer(jmsSession, settings, 0)
jmsProducer.createMessage(
JmsMapMessage(
- Map(
+ Map[String, Any](
"string" -> "string",
"int" -> 1,
"boolean" -> true,
diff --git
a/kinesis/src/test/scala/org/apache/pekko/stream/connectors/kinesis/KinesisFlowSpec.scala
b/kinesis/src/test/scala/org/apache/pekko/stream/connectors/kinesis/KinesisFlowSpec.scala
index 4eb0f8722..eaf9a1074 100644
---
a/kinesis/src/test/scala/org/apache/pekko/stream/connectors/kinesis/KinesisFlowSpec.scala
+++
b/kinesis/src/test/scala/org/apache/pekko/stream/connectors/kinesis/KinesisFlowSpec.scala
@@ -32,7 +32,6 @@ import org.scalatest.matchers.should.Matchers
import software.amazon.awssdk.core.SdkBytes
import software.amazon.awssdk.services.kinesis.model._
-import scala.annotation.nowarn
import scala.jdk.CollectionConverters._
class KinesisFlowSpec extends AnyWordSpec with Matchers with KinesisMock with
LogCapturing {
@@ -111,10 +110,9 @@ class KinesisFlowSpec extends AnyWordSpec with Matchers
with KinesisMock with Lo
.run()
}
- @nowarn("msg=deprecated")
trait KinesisFlowWithContextProbe { self: Settings =>
val streamName = "stream-name"
- val recordStream = Stream
+ val recordStream = LazyList
.from(1)
.map(i =>
(PutRecordsRequestEntry
@@ -123,7 +121,7 @@ class KinesisFlowSpec extends AnyWordSpec with Matchers
with KinesisMock with Lo
.data(SdkBytes.fromByteBuffer(ByteString(i).asByteBuffer))
.build(),
i))
- val resultStream = Stream
+ val resultStream = LazyList
.from(1)
.map(i => (PutRecordsResultEntry.builder().build(), i))
diff --git
a/kinesis/src/test/scala/org/apache/pekko/stream/connectors/kinesis/Valve.scala
b/kinesis/src/test/scala/org/apache/pekko/stream/connectors/kinesis/Valve.scala
index c9d039ce9..a1574115c 100644
---
a/kinesis/src/test/scala/org/apache/pekko/stream/connectors/kinesis/Valve.scala
+++
b/kinesis/src/test/scala/org/apache/pekko/stream/connectors/kinesis/Valve.scala
@@ -18,6 +18,7 @@ import pekko.stream.stage.{ GraphStageLogic,
GraphStageWithMaterializedValue, In
import pekko.stream._
import SwitchMode.{ Close, Open }
+import scala.annotation.nowarn
import scala.concurrent.{ Future, Promise }
/**
@@ -86,6 +87,7 @@ final class Valve[A](mode: SwitchMode) extends
GraphStageWithMaterializedValue[F
(logic, logic.promise.future)
}
+ @nowarn("msg=inferred structural type")
private class ValveGraphStageLogic(shape: Shape, var mode: SwitchMode)
extends GraphStageLogic(shape)
with InHandler
@@ -93,7 +95,7 @@ final class Valve[A](mode: SwitchMode) extends
GraphStageWithMaterializedValue[F
val promise = Promise[ValveSwitch]()
- private val switch = new ValveSwitch {
+ private val switch: ValveSwitch = new ValveSwitch {
val flipCallback = getAsyncCallback[(SwitchMode, Promise[Boolean])] {
case (flipToMode, promise) =>
diff --git
a/mqtt-streaming/src/main/scala/org/apache/pekko/stream/connectors/mqtt/streaming/impl/MqttFrameStage.scala
b/mqtt-streaming/src/main/scala/org/apache/pekko/stream/connectors/mqtt/streaming/impl/MqttFrameStage.scala
index 8ae10093c..d7e05cfce 100644
---
a/mqtt-streaming/src/main/scala/org/apache/pekko/stream/connectors/mqtt/streaming/impl/MqttFrameStage.scala
+++
b/mqtt-streaming/src/main/scala/org/apache/pekko/stream/connectors/mqtt/streaming/impl/MqttFrameStage.scala
@@ -47,7 +47,7 @@ import scala.collection.immutable
} else {
Left(new IllegalStateException(s"Max packet size of $maxPacketSize
exceeded with $packetSize"))
}
- case _: Left[BufferUnderflow.type, Int] @unchecked =>
+ case Left(_) =>
Right((bytesToEmit, bytesReceived))
}
}
diff --git
a/mqtt-streaming/src/main/scala/org/apache/pekko/stream/connectors/mqtt/streaming/scaladsl/MqttSession.scala
b/mqtt-streaming/src/main/scala/org/apache/pekko/stream/connectors/mqtt/streaming/scaladsl/MqttSession.scala
index f5da6ad53..bd06d2abf 100644
---
a/mqtt-streaming/src/main/scala/org/apache/pekko/stream/connectors/mqtt/streaming/scaladsl/MqttSession.scala
+++
b/mqtt-streaming/src/main/scala/org/apache/pekko/stream/connectors/mqtt/streaming/scaladsl/MqttSession.scala
@@ -26,6 +26,7 @@ import pekko.stream.scaladsl.{ BroadcastHub, Flow, Keep,
Source }
import pekko.util.ByteString
import pekko.{ Done, NotUsed }
+import scala.annotation.nowarn
import scala.concurrent.{ Future, Promise }
import scala.util.control.NoStackTrace
import scala.util.{ Failure, Success }
@@ -300,6 +301,7 @@ final class ActorMqttClientSession(settings:
MqttSessionSettings)(implicit syste
}
.mapMaterializedValue(_ => NotUsed)
+ @nowarn("msg=exhaustive")
private[streaming] override def eventFlow[A](connectionId: ByteString):
EventFlow[A] =
Flow[ByteString]
.watch(clientConnector.toClassic)
@@ -654,6 +656,7 @@ final class ActorMqttServerSession(settings:
MqttSessionSettings)(implicit syste
}
.mapMaterializedValue(_ => NotUsed)
+ @nowarn("msg=exhaustive")
override def eventFlow[A](connectionId: ByteString): EventFlow[A] =
Flow[ByteString]
.watch(serverConnector.toClassic)
diff --git
a/s3/src/test/scala/org/apache/pekko/stream/connectors/s3/scaladsl/Generators.scala
b/s3/src/test/scala/org/apache/pekko/stream/connectors/s3/scaladsl/Generators.scala
index 509b0d599..3adb5dda2 100644
---
a/s3/src/test/scala/org/apache/pekko/stream/connectors/s3/scaladsl/Generators.scala
+++
b/s3/src/test/scala/org/apache/pekko/stream/connectors/s3/scaladsl/Generators.scala
@@ -20,7 +20,6 @@ package org.apache.pekko.stream.connectors.s3.scaladsl
import org.scalacheck.Gen
import scala.annotation.nowarn
-import scala.language.postfixOps
object Generators {
val MaxBucketLength: Int = 63
diff --git
a/s3/src/test/scala/org/apache/pekko/stream/connectors/s3/scaladsl/S3ExtSpec.scala
b/s3/src/test/scala/org/apache/pekko/stream/connectors/s3/scaladsl/S3ExtSpec.scala
index b1084af65..ddddfa998 100644
---
a/s3/src/test/scala/org/apache/pekko/stream/connectors/s3/scaladsl/S3ExtSpec.scala
+++
b/s3/src/test/scala/org/apache/pekko/stream/connectors/s3/scaladsl/S3ExtSpec.scala
@@ -26,7 +26,7 @@ import scala.jdk.CollectionConverters._
class S3ExtSpec extends AnyFlatSpecLike with Matchers {
it should "reuse application config from actor system" in {
val config = ConfigFactory.parseMap(
- Map(
+ Map[String, Any](
"pekko.connectors.s3.endpoint-url" -> "http://localhost:8001",
"pekko.connectors.s3.path-style-access" -> true).asJava)
implicit val system: ActorSystem = ActorSystem.create("s3", config)
diff --git
a/s3/src/test/scala/org/apache/pekko/stream/connectors/s3/scaladsl/S3IntegrationSpec.scala
b/s3/src/test/scala/org/apache/pekko/stream/connectors/s3/scaladsl/S3IntegrationSpec.scala
index 90893a5d5..72e30533b 100644
---
a/s3/src/test/scala/org/apache/pekko/stream/connectors/s3/scaladsl/S3IntegrationSpec.scala
+++
b/s3/src/test/scala/org/apache/pekko/stream/connectors/s3/scaladsl/S3IntegrationSpec.scala
@@ -748,7 +748,7 @@ trait S3IntegrationSpec
_ <- pekko.pattern.after(25.seconds)(Future {
sharedKillSwitch.abort(AbortException)
})
- _ <- multiPartUpload.recover {
+ _ <- multiPartUpload.recover[Any] {
case AbortException => ()
}
incomplete <- S3.listMultipartUpload(defaultBucket,
None).withAttributes(attributes).runWith(Sink.seq)
@@ -801,7 +801,7 @@ trait S3IntegrationSpec
_ <- pekko.pattern.after(25.seconds)(Future {
sharedKillSwitch.abort(AbortException)
})
- _ <- multiPartUpload.recover {
+ _ <- multiPartUpload.recover[Any] {
case AbortException => ()
}
incomplete <- S3.listMultipartUpload(defaultBucket,
None).withAttributes(attributes).runWith(Sink.seq)
@@ -856,7 +856,7 @@ trait S3IntegrationSpec
_ <- pekko.pattern.after(25.seconds)(Future {
sharedKillSwitch.abort(AbortException)
})
- _ <- multiPartUpload.recover {
+ _ <- multiPartUpload.recover[Any] {
case AbortException => ()
}
incomplete <- S3.listMultipartUpload(defaultBucket,
None).withAttributes(attributes).runWith(Sink.seq)
diff --git a/solr/src/test/scala/docs/scaladsl/SolrSpec.scala
b/solr/src/test/scala/docs/scaladsl/SolrSpec.scala
index 9b7bec68c..27a75179b 100644
--- a/solr/src/test/scala/docs/scaladsl/SolrSpec.scala
+++ b/solr/src/test/scala/docs/scaladsl/SolrSpec.scala
@@ -36,6 +36,7 @@ import org.apache.solr.common.SolrInputDocument
import org.scalatest.concurrent.ScalaFutures
import org.scalatest.BeforeAndAfterAll
+import scala.annotation.nowarn
import scala.collection.immutable
import scala.concurrent.duration._
import scala.concurrent.{ ExecutionContext, Future }
@@ -719,6 +720,7 @@ class SolrSpec extends AnyWordSpec with Matchers with
BeforeAndAfterAll with Sca
TestKit.shutdownActorSystem(system)
}
+ @nowarn("msg=deprecated")
private def setupCluster(): Unit = {
val targetDir = new File("solr/target")
val testWorkingDir =
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]