diff --git a/project/Common.scala b/project/Common.scala index 262f88f..3f667c2 100644 --- a/project/Common.scala +++ b/project/Common.scala @@ -16,19 +16,14 @@ object Common extends AutoPlugin { organizationHomepage := Some(url("http://www.tmt.org")), resolvers += "jitpack" at "https://jitpack.io", + // Scala 3 (3.6.4). The following Scala 2-only flags were removed in the port: + // -Wdead-code, -Xlint:_,-missing-interpolator, -Xsource:3, -Xcheckinit, -Xasync scalacOptions ++= Seq( "-encoding", "UTF-8", "-feature", "-unchecked", - "-deprecation", - //-W Options - "-Wdead-code", - //-X Options - "-Xlint:_,-missing-interpolator", - "-Xsource:3", - "-Xcheckinit", - "-Xasync" + "-deprecation" ), Compile / doc / javacOptions ++= Seq("-Xdoclint:none"), Test / testOptions ++= Seq( @@ -47,4 +42,4 @@ object Common extends AutoPlugin { // case Some("false") => false // case _ => true // } -} +} \ No newline at end of file diff --git a/project/Dependencies.scala b/project/Dependencies.scala index 1f7cfa0..837858d 100644 --- a/project/Dependencies.scala +++ b/project/Dependencies.scala @@ -1,47 +1,34 @@ import sbt._ -import org.portablescala.sbtplatformdeps.PlatformDepsPlugin.autoImport._ object Dependencies { val vbdsServer = Seq( CSW.`csw-location-client`, CSW.`csw-network-utils`, - Akka.`akka-stream-typed`, - Akka.`akka-actor-typed`, - Akka.`akka-cluster-typed`, - Akka.`akka-slf4j`, - AkkaHttp.`akka-http`, - AkkaHttp.`akka-http-spray-json`, + Pekko.`pekko-stream-typed`, + Pekko.`pekko-actor-typed`, + Pekko.`pekko-cluster-typed`, + Pekko.`pekko-slf4j`, + PekkoHttp.`pekko-http`, + PekkoHttp.`pekko-http-spray-json`, Libs.`scopt`, Libs.`boopickle`, -// Chill.`chill-akka`, Libs.`logback-classic`, Libs.`commons-io` % Test, Libs.`scalatest` % Test, - AkkaHttp.`akka-http-testkit` % Test, - Akka.`akka-multi-node-testkit` % Test + PekkoHttp.`pekko-http-testkit` % Test, + Pekko.`pekko-multi-node-testkit` % Test ) val vbdsClient = Seq( CSW.`csw-location-client`, - Akka.`akka-stream`, - Akka.`akka-slf4j`, - AkkaHttp.`akka-http`, + Pekko.`pekko-stream`, + Pekko.`pekko-slf4j`, + PekkoHttp.`pekko-http`, Libs.`scopt`, Libs.`logback-classic`, Libs.`scalatest` % Test, - AkkaHttp.`akka-http-testkit` % Test + PekkoHttp.`pekko-http-testkit` % Test ) - - // ScalaJS web client scala dependencies - val webClient = Def.setting(Seq( - "org.scala-js" %%% "scalajs-dom" % "1.1.0", - "com.lihaoyi" %%% "scalatags" % "0.9.4", - "com.github.japgolly.scalacss" %%% "core" % "0.7.0", - "com.github.japgolly.scalacss" %%% "ext-scalatags" % "0.7.0", - "com.lihaoyi" %%% "upickle" % "1.4.0", - "org.scalatest" %%% "scalatest" % "3.2.9" % "test" - )) - -} +} \ No newline at end of file diff --git a/project/Libs.scala b/project/Libs.scala index 50a04ae..8c0019d 100644 --- a/project/Libs.scala +++ b/project/Libs.scala @@ -1,54 +1,52 @@ import sbt._ object Libs { - val ScalaVersion = "2.13.8" - - val `scalatest` = "org.scalatest" %% "scalatest" % "3.2.11" //Apache License 2.0 - val `scala-async` = "org.scala-lang.modules" %% "scala-async" % "1.0.1" //BSD 3-clause "New" or "Revised" License - val `scopt` = "com.github.scopt" %% "scopt" % "4.0.1" //MIT License - val `logback-classic` = "ch.qos.logback" % "logback-classic" % "1.2.10" // GNU Lesser General Public License version 2.1 - val `akka-management-cluster-http` = "com.lightbend.akka" %% "akka-management-cluster-http" % "1.1.3" //N/A at the moment - val `boopickle` = "io.suzaku" %% "boopickle" % "1.4.0" // Apache License 2.0 - val `commons-io` = "commons-io" % "commons-io" % "2.11.0" // Apache License 2.0 + val ScalaVersion = "3.6.4" + + val `scalatest` = "org.scalatest" %% "scalatest" % "3.2.19" // Apache License 2.0 (aligned w/ csw-6.0.0) + val `scopt` = "com.github.scopt" %% "scopt" % "4.1.0" // MIT License (aligned w/ csw-6.0.0) + val `logback-classic` = "ch.qos.logback" % "logback-classic" % "1.5.17" // EPL v1.0 / LGPL 2.1 (aligned w/ csw-6.0.0) + val `boopickle` = "io.suzaku" %% "boopickle" % "1.4.0" // Apache License 2.0 (Scala 3-capable) + val `commons-io` = "commons-io" % "commons-io" % "2.11.0" // Apache License 2.0 + + // Removed in the csw-6.0.0 / Pekko / Scala 3 port: + // scala-async -- Scala 2-only, and was not wired into either module + // akka-management-cluster-http -- defined-but-unused } object CSW { - private val Org = "com.github.tmtsoftware.csw" - // private val Version = "0.1.0-SNAPSHOT" - private val Version = "5.0.0" + private val Org = "com.github.tmtsoftware.csw" + private val Version = "6.0.0" - val `csw-network-utils` = Org %% "csw-network-utils" % Version + val `csw-network-utils` = Org %% "csw-network-utils" % Version val `csw-location-client` = Org %% "csw-location-client" % Version - val `csw-commons` = Org %% "csw-commons" % Version + val `csw-commons` = Org %% "csw-commons" % Version } -object Akka { - val Version = "2.6.18" //all akka is Apache License 2.0 - val `akka-stream` = "com.typesafe.akka" %% "akka-stream" % Version - val `akka-stream-typed` = "com.typesafe.akka" %% "akka-stream-typed" % Version - val `akka-stream-testkit` = "com.typesafe.akka" %% "akka-stream-testkit" % Version - val `akka-actor` = "com.typesafe.akka" %% "akka-actor" % Version - val `akka-actor-typed` = "com.typesafe.akka" %% "akka-actor-typed" % Version - val `akka-testkit-typed` = "com.typesafe.akka" %% "akka-testkit-typed" % Version - val `akka-distributed-data` = "com.typesafe.akka" %% "akka-distributed-data" % Version - val `akka-multi-node-testkit` = "com.typesafe.akka" %% "akka-multi-node-testkit" % Version - val `akka-cluster` = "com.typesafe.akka" %% "akka-cluster" % Version - val `akka-cluster-tools` = "com.typesafe.akka" %% "akka-cluster-tools" % Version - val `akka-cluster-typed` = "com.typesafe.akka" %% "akka-cluster-typed" % Version - val `akka-slf4j` = "com.typesafe.akka" %% "akka-slf4j" % Version +object Pekko { // all pekko is Apache License 2.0 + val Version = "1.1.3" + val Org = "org.apache.pekko" + + val `pekko-stream` = Org %% "pekko-stream" % Version + val `pekko-stream-typed` = Org %% "pekko-stream-typed" % Version + val `pekko-stream-testkit` = Org %% "pekko-stream-testkit" % Version + val `pekko-actor` = Org %% "pekko-actor" % Version + val `pekko-actor-typed` = Org %% "pekko-actor-typed" % Version + val `pekko-actor-testkit-typed`= Org %% "pekko-actor-testkit-typed" % Version + val `pekko-distributed-data` = Org %% "pekko-distributed-data" % Version + val `pekko-multi-node-testkit` = Org %% "pekko-multi-node-testkit" % Version + val `pekko-cluster` = Org %% "pekko-cluster" % Version + val `pekko-cluster-tools` = Org %% "pekko-cluster-tools" % Version + val `pekko-cluster-typed` = Org %% "pekko-cluster-typed" % Version + val `pekko-slf4j` = Org %% "pekko-slf4j" % Version } -object AkkaHttp { //ApacheV2 - val Version = "10.2.7" - val `akka-http` = "com.typesafe.akka" %% "akka-http" % Version - val `akka-http-core` = "com.typesafe.akka" %% "akka-http-core" % Version - val `akka-http-testkit` = "com.typesafe.akka" %% "akka-http-testkit" % Version - val `akka-http-spray-json` = "com.typesafe.akka" %% "akka-http-spray-json" % Version -} - -//object Chill { -// val Version = "0.10.0" -// val `chill-akka` = "com.twitter" %% "chill-akka" % Version //Apache License 2.0 -//} - +object PekkoHttp { // Apache License 2.0 + val Version = "1.1.0" + val Org = "org.apache.pekko" + val `pekko-http` = Org %% "pekko-http" % Version + val `pekko-http-core` = Org %% "pekko-http-core" % Version + val `pekko-http-testkit` = Org %% "pekko-http-testkit" % Version + val `pekko-http-spray-json` = Org %% "pekko-http-spray-json" % Version +} \ No newline at end of file diff --git a/project/build.properties b/project/build.properties index 19479ba..dabdb15 100644 --- a/project/build.properties +++ b/project/build.properties @@ -1 +1 @@ -sbt.version=1.5.2 +sbt.version=1.12.11 diff --git a/project/plugins.sbt b/project/plugins.sbt index d35fdd2..ee0a246 100755 --- a/project/plugins.sbt +++ b/project/plugins.sbt @@ -1,16 +1,7 @@ //addSbtPlugin("org.scalastyle" %% "scalastyle-sbt-plugin" % "1.0.0") //addSbtPlugin("com.geirsson" % "sbt-scalafmt" % "1.4.0") //addSbtPlugin("org.scoverage" % "sbt-scoverage" % "1.9.3") -addSbtPlugin("com.github.sbt" % "sbt-native-packager" % "1.9.16") -addSbtPlugin("com.eed3si9n" % "sbt-buildinfo" % "0.10.0") -addSbtPlugin("com.typesafe.sbt" % "sbt-multi-jvm" % "0.4.0") -addSbtPlugin("org.portable-scala" % "sbt-scalajs-crossproject" % "1.1.0") - -// web client -addSbtPlugin("org.scala-js" % "sbt-scalajs" % "1.8.0") - -// Requires local plugin build and publishLocal, since existing plugin was abandoned -//addSbtPlugin("com.lihaoyi" % "workbench" % "0.4.2") - -addSbtPlugin("org.scala-js" % "sbt-jsdependencies" % "1.0.2") -addSbtPlugin("com.timushev.sbt" % "sbt-updates" % "0.6.1") +addSbtPlugin("com.github.sbt" % "sbt-native-packager" % "1.9.16") +addSbtPlugin("com.eed3si9n" % "sbt-buildinfo" % "0.10.0") +addSbtPlugin("com.typesafe.sbt" % "sbt-multi-jvm" % "0.4.0") +addSbtPlugin("com.timushev.sbt" % "sbt-updates" % "0.6.1") diff --git a/vbds-client/src/main/scala/vbds/client/FileUploader.scala b/vbds-client/src/main/scala/vbds/client/FileUploader.scala index 8ede0ee..f3f7b7e 100644 --- a/vbds-client/src/main/scala/vbds/client/FileUploader.scala +++ b/vbds-client/src/main/scala/vbds/client/FileUploader.scala @@ -2,14 +2,14 @@ package vbds.client import java.nio.file.Path -import akka.{Done, NotUsed} -import akka.actor.ActorSystem -import akka.http.scaladsl.Http -import akka.http.scaladsl.marshalling.Marshal -import akka.http.scaladsl.model.Multipart.FormData -import akka.http.scaladsl.model._ -import akka.stream.ThrottleMode -import akka.stream.scaladsl._ +import org.apache.pekko.{Done, NotUsed} +import org.apache.pekko.actor.ActorSystem +import org.apache.pekko.http.scaladsl.Http +import org.apache.pekko.http.scaladsl.marshalling.Marshal +import org.apache.pekko.http.scaladsl.model.Multipart.FormData +import org.apache.pekko.http.scaladsl.model._ +import org.apache.pekko.stream.ThrottleMode +import org.apache.pekko.stream.scaladsl._ import scala.concurrent.Future import scala.concurrent.duration._ @@ -55,4 +55,4 @@ class FileUploader(chunkSize: Int = 1024 * 1024)(implicit val system: ActorSyste upload(source.throttle(1, delay, 1, ThrottleMode.Shaping)) else upload(source) } -} +} \ No newline at end of file diff --git a/vbds-client/src/main/scala/vbds/client/VbdsClient.scala b/vbds-client/src/main/scala/vbds/client/VbdsClient.scala index fe653ce..379d02b 100644 --- a/vbds-client/src/main/scala/vbds/client/VbdsClient.scala +++ b/vbds-client/src/main/scala/vbds/client/VbdsClient.scala @@ -3,17 +3,16 @@ package vbds.client import java.io.File import java.nio.file.Path import java.time.Instant - -import akka.Done -import akka.actor.{ActorRef, ActorSystem} -import akka.event.{LogSource, Logging} -import akka.http.scaladsl.Http -import akka.http.scaladsl.model.* -import akka.http.scaladsl.model.ws.Message -import akka.stream.scaladsl.MergeHub +import org.apache.pekko.Done +import org.apache.pekko.actor.{ActorRef, ActorSystem} +import org.apache.pekko.event.{LogSource, Logging} +import org.apache.pekko.http.scaladsl.Http +import org.apache.pekko.http.scaladsl.model.* +import org.apache.pekko.http.scaladsl.model.ws.Message +import org.apache.pekko.stream.scaladsl.MergeHub import vbds.client.VbdsClient.Subscription -import scala.concurrent.{Await, Future} +import scala.concurrent.{Await, ExecutionContextExecutor, Future} import scala.concurrent.duration.* import scala.util.{Failure, Success, Try} @@ -37,7 +36,7 @@ object VbdsClient { //noinspection HttpUrlsUsage class VbdsClient(name: String, host: String, port: Int, chunkSize: Int = 1024 * 1024)(implicit val system: ActorSystem) { - implicit val executionContext = system.dispatcher + implicit val executionContext: ExecutionContextExecutor = system.dispatcher val adminRoute = "/vbds/admin/streams" val accessRoute = "/vbds/access/streams" val transferRoute = "/vbds/transfer/streams" @@ -197,4 +196,4 @@ class VbdsClient(name: String, host: String, port: Int, chunkSize: Int = 1024 * wsListener.subscribe(uri, receiver, outSource) } -} +} \ No newline at end of file diff --git a/vbds-client/src/main/scala/vbds/client/WebSocketActor.scala b/vbds-client/src/main/scala/vbds/client/WebSocketActor.scala index 621f23b..0f5ad7b 100644 --- a/vbds-client/src/main/scala/vbds/client/WebSocketActor.scala +++ b/vbds-client/src/main/scala/vbds/client/WebSocketActor.scala @@ -3,16 +3,17 @@ package vbds.client import java.io.{File, FileOutputStream} import java.nio.file.Path -import akka.NotUsed -import akka.actor.{Actor, ActorLogging, ActorRef, ActorSystem, Props} -import akka.http.scaladsl.model.ws.{BinaryMessage, Message, TextMessage} -import akka.util.{ByteString, Timeout} -import akka.stream.scaladsl.{Sink, Source} +import org.apache.pekko.NotUsed +import org.apache.pekko.actor.{Actor, ActorLogging, ActorRef, ActorSystem, Props} +import org.apache.pekko.http.scaladsl.model.ws.{BinaryMessage, Message, TextMessage} +import org.apache.pekko.util.{ByteString, Timeout} +import org.apache.pekko.stream.scaladsl.{Sink, Source} import scala.concurrent.Future import scala.concurrent.duration._ import scala.util._ -import akka.pattern.ask +import scala.compiletime.uninitialized +import org.apache.pekko.pattern.ask /** * An actor that receives websocket messages from the VBDS server. @@ -74,9 +75,9 @@ class WebSocketActor(name: String, import system.dispatcher var count = 0 - var file: File = _ - var os: FileOutputStream = _ - implicit val askTimeout = Timeout(20.seconds) + var file: File = uninitialized + var os: FileOutputStream = uninitialized + implicit val askTimeout: Timeout = Timeout(20.seconds) log.debug(s"$name: Started WebSocketActor") @@ -142,4 +143,4 @@ class WebSocketActor(name: String, } } -} +} \ No newline at end of file diff --git a/vbds-client/src/main/scala/vbds/client/WebSocketListener.scala b/vbds-client/src/main/scala/vbds/client/WebSocketListener.scala index 254163f..cce08df 100644 --- a/vbds-client/src/main/scala/vbds/client/WebSocketListener.scala +++ b/vbds-client/src/main/scala/vbds/client/WebSocketListener.scala @@ -1,11 +1,11 @@ package vbds.client -import akka.NotUsed -import akka.actor.{ActorRef, ActorSystem} -import akka.http.scaladsl.Http -import akka.http.scaladsl.model.ws.{Message, WebSocketRequest} -import akka.http.scaladsl.model.{HttpResponse, Uri} -import akka.stream.scaladsl.{Flow, Sink, Source} +import org.apache.pekko.NotUsed +import org.apache.pekko.actor.{ActorRef, ActorSystem} +import org.apache.pekko.http.scaladsl.Http +import org.apache.pekko.http.scaladsl.model.ws.{Message, WebSocketRequest} +import org.apache.pekko.http.scaladsl.model.{HttpResponse, Uri} +import org.apache.pekko.stream.scaladsl.{Flow, Sink, Source} import vbds.client.WebSocketListener.SubscribeResult import scala.concurrent.{Future, Promise} @@ -61,4 +61,4 @@ class WebSocketListener(implicit val system: ActorSystem) { new SubscribeResult(upgradeResponse.map(_.response), promise) } -} +} \ No newline at end of file diff --git a/vbds-client/src/main/scala/vbds/client/app/VbdsClientApp.scala b/vbds-client/src/main/scala/vbds/client/app/VbdsClientApp.scala index 7b893d9..0bb10c7 100644 --- a/vbds-client/src/main/scala/vbds/client/app/VbdsClientApp.scala +++ b/vbds-client/src/main/scala/vbds/client/app/VbdsClientApp.scala @@ -1,15 +1,15 @@ package vbds.client.app import java.io.File -import akka.Done -import akka.actor.{Actor, ActorLogging, ActorSystem, Props} -import akka.http.scaladsl.model.HttpResponse +import org.apache.pekko.Done +import org.apache.pekko.actor.{Actor, ActorLogging, ActorSystem, Props} +import org.apache.pekko.http.scaladsl.model.HttpResponse import csw.location.api.models.{ComponentId, ComponentType, HttpLocation} import csw.location.api.models.Connection.HttpConnection import csw.location.client.scaladsl.HttpLocationServiceFactory import csw.prefix.models.{Prefix, Subsystem} import vbds.client.VbdsClient -import akka.actor.typed.scaladsl.adapter.ClassicActorSystemOps +import org.apache.pekko.actor.typed.scaladsl.adapter.ClassicActorSystemOps import scala.concurrent.{Await, Future} import scala.concurrent.duration.* @@ -19,8 +19,8 @@ import vbds.server.app.BuildInfo /** * A VIZ Bulk Data System HTTP client command line application. */ -object VbdsClientApp extends App { - implicit val system = ActorSystem("vbdsClient") +object VbdsClientApp { + implicit val system: ActorSystem = ActorSystem("vbdsClient") // Command line options private case class Options(name: String = "vbds", @@ -118,16 +118,18 @@ object VbdsClientApp extends App { } // Parse the command line options - parser.parse(args, Options()) match { - case Some(options) => - try { - run(options) - } catch { - case e: Throwable => - e.printStackTrace() - System.exit(1) - } - case None => System.exit(1) + def main(args: Array[String]): Unit = { + parser.parse(args, Options()) match { + case Some(options) => + try { + run(options) + } catch { + case e: Throwable => + e.printStackTrace() + System.exit(1) + } + case None => System.exit(1) + } } // Run the application (The actor system is only used locally, no need for remote) @@ -224,4 +226,4 @@ object VbdsClientApp extends App { System.exit(1) } } -} +} \ No newline at end of file diff --git a/vbds-client/src/test/scala/vbds/client/StreamTest.scala b/vbds-client/src/test/scala/vbds/client/StreamTest.scala index 784ef40..3f7f10f 100644 --- a/vbds-client/src/test/scala/vbds/client/StreamTest.scala +++ b/vbds-client/src/test/scala/vbds/client/StreamTest.scala @@ -1,12 +1,12 @@ package vbds.client -import akka.actor.ActorSystem -import akka.stream.scaladsl.{BroadcastHub, Keep, Sink, Source} +import org.apache.pekko.actor.ActorSystem +import org.apache.pekko.stream.scaladsl.{BroadcastHub, Keep, Sink, Source} import scala.concurrent.Future import scala.util.{Failure, Success} object StreamTest extends App { - implicit val system = ActorSystem("StreamTest") + implicit val system: ActorSystem = ActorSystem("StreamTest") import system._ val source = Source(0 to 5) @@ -36,4 +36,4 @@ object StreamTest extends App { println(s"Error: $ex") system.terminate() } -} +} \ No newline at end of file diff --git a/vbds-server/src/main/resources/application.conf b/vbds-server/src/main/resources/application.conf index 36ac69a..a4a260d 100644 --- a/vbds-server/src/main/resources/application.conf +++ b/vbds-server/src/main/resources/application.conf @@ -1,7 +1,7 @@ include required("logging.conf") vbds-system { - akka.actor.provider = cluster + pekko.actor.provider = cluster } csw-logging { @@ -10,11 +10,11 @@ csw-logging { logLevel = debug // Log level for slf4j messages slf4jLogLevel = info - // Log level for Akka messages, should be >= akka.loglevel - akkaLogLevel = debug + // Log level for Pekko messages, should be >= pekko.loglevel + pekkoLogLevel = debug } -akka { +pekko { loglevel = DEBUG log-dead-letters-during-shutdown = off diff --git a/vbds-server/src/main/scala/vbds/server/actors/AccessApi.scala b/vbds-server/src/main/scala/vbds/server/actors/AccessApi.scala index f5d5f3c..42113ff 100644 --- a/vbds-server/src/main/scala/vbds/server/actors/AccessApi.scala +++ b/vbds-server/src/main/scala/vbds/server/actors/AccessApi.scala @@ -1,13 +1,13 @@ package vbds.server.actors -import akka.NotUsed -import akka.actor.typed.{ActorRef, Scheduler} -import akka.stream.scaladsl.Sink -import akka.util.ByteString -import akka.util.Timeout +import org.apache.pekko.NotUsed +import org.apache.pekko.actor.typed.{ActorRef, Scheduler} +import org.apache.pekko.stream.scaladsl.Sink +import org.apache.pekko.util.ByteString +import org.apache.pekko.util.Timeout import vbds.server.actors.SharedDataActor.SharedDataActorMessages import vbds.server.models.AccessInfo -import akka.actor.typed.scaladsl.AskPattern._ +import org.apache.pekko.actor.typed.scaladsl.AskPattern._ import vbds.server.routes.AccessRoute.WebsocketResponseActorMsg import scala.concurrent.{ExecutionContext, Future} @@ -44,7 +44,7 @@ class AccessApiImpl(sharedDataActor: ActorRef[SharedDataActorMessages])(implicit } def listSubscriptions(): Future[Set[AccessInfo]] = { - sharedDataActor.ask(ListSubscriptions) + sharedDataActor.ask(ListSubscriptions.apply) } def subscriptionExists(id: String): Future[Boolean] = { diff --git a/vbds-server/src/main/scala/vbds/server/actors/AdminApi.scala b/vbds-server/src/main/scala/vbds/server/actors/AdminApi.scala index 03df5e3..f586292 100644 --- a/vbds-server/src/main/scala/vbds/server/actors/AdminApi.scala +++ b/vbds-server/src/main/scala/vbds/server/actors/AdminApi.scala @@ -1,13 +1,13 @@ package vbds.server.actors -import akka.actor.typed.{ActorRef, Scheduler} -import akka.util.Timeout +import org.apache.pekko.actor.typed.{ActorRef, Scheduler} +import org.apache.pekko.util.Timeout import vbds.server.actors.SharedDataActor.SharedDataActorMessages import vbds.server.models.StreamInfo import scala.concurrent.{ExecutionContext, Future} import scala.concurrent.duration._ -import akka.actor.typed.scaladsl.AskPattern._ +import org.apache.pekko.actor.typed.scaladsl.AskPattern._ /** * Internal Admin API @@ -32,7 +32,7 @@ class AdminApiImpl(sharedDataActor: ActorRef[SharedDataActorMessages])(implicit import SharedDataActor._ def listStreams(): Future[Set[StreamInfo]] = { - sharedDataActor.ask(ListStreams) + sharedDataActor.ask(ListStreams.apply) } def streamExists(name: String): Future[Boolean] = { diff --git a/vbds-server/src/main/scala/vbds/server/actors/AkkaTypedExtension.scala b/vbds-server/src/main/scala/vbds/server/actors/AkkaTypedExtension.scala index 78629b1..54277ec 100644 --- a/vbds-server/src/main/scala/vbds/server/actors/AkkaTypedExtension.scala +++ b/vbds-server/src/main/scala/vbds/server/actors/AkkaTypedExtension.scala @@ -1,10 +1,10 @@ package vbds.server.actors -import akka.actor.typed.Scheduler -import akka.actor.typed.SpawnProtocol.Spawn -import akka.actor.typed._ -import akka.actor.typed.scaladsl.AskPattern._ -import akka.util.Timeout +import org.apache.pekko.actor.typed.Scheduler +import org.apache.pekko.actor.typed.SpawnProtocol.Spawn +import org.apache.pekko.actor.typed._ +import org.apache.pekko.actor.typed.scaladsl.AskPattern._ +import org.apache.pekko.util.Timeout import scala.concurrent.Await import scala.concurrent.duration.{DurationInt, FiniteDuration} diff --git a/vbds-server/src/main/scala/vbds/server/actors/SharedDataActor.scala b/vbds-server/src/main/scala/vbds/server/actors/SharedDataActor.scala index 78c41ab..f7f7874 100644 --- a/vbds-server/src/main/scala/vbds/server/actors/SharedDataActor.scala +++ b/vbds-server/src/main/scala/vbds/server/actors/SharedDataActor.scala @@ -1,29 +1,30 @@ package vbds.server.actors import java.net.InetSocketAddress -import akka.{Done, NotUsed} -import akka.actor.typed.scaladsl.{AbstractBehavior, ActorContext, Behaviors} -import akka.cluster.ddata.typed.scaladsl.{DistributedData, ReplicatorMessageAdapter} -import akka.cluster.ddata.typed.scaladsl.Replicator.* -import akka.stream.ClosedShape -import akka.stream.scaladsl.{Broadcast, Flow, GraphDSL, Keep, Merge, RunnableGraph, Sink, Source} -import akka.util.{ByteString, Timeout} +import org.apache.pekko.{Done, NotUsed} +import org.apache.pekko.actor.typed.scaladsl.{AbstractBehavior, ActorContext, Behaviors} +import org.apache.pekko.cluster.ddata.typed.scaladsl.{DistributedData, ReplicatorMessageAdapter} +import org.apache.pekko.cluster.ddata.typed.scaladsl.Replicator.* +import org.apache.pekko.stream.ClosedShape +import org.apache.pekko.stream.scaladsl.{Broadcast, Flow, GraphDSL, Keep, Merge, RunnableGraph, Sink, Source} +import org.apache.pekko.util.{ByteString, Timeout} import vbds.server.models.{AccessInfo, ServerInfo, StreamInfo} -import akka.http.scaladsl.Http -import akka.http.scaladsl.model.{ContentTypes, HttpEntity, HttpMethods, HttpRequest, HttpResponse} +import org.apache.pekko.http.scaladsl.Http +import org.apache.pekko.http.scaladsl.model.{ContentTypes, HttpEntity, HttpMethods, HttpRequest, HttpResponse} import scala.concurrent.duration.* import scala.concurrent.{ExecutionContext, Future} import scala.util.{Failure, Success} +import scala.compiletime.uninitialized import SharedDataActor.* -import akka.actor.typed.{ActorRef, ActorSystem, Behavior, PostStop, Scheduler, Signal, SpawnProtocol} -import akka.cluster.ddata.{ORSet, ORSetKey, SelfUniqueAddress} -import akka.actor.typed.scaladsl.AskPattern.* -import akka.cluster.ClusterEvent.{MemberEvent, ReachabilityEvent} -import akka.cluster.typed.{Cluster, Leave} -import akka.event.Logging.{DebugLevel, ErrorLevel, InfoLevel, LogLevel, WarningLevel} +import org.apache.pekko.actor.typed.{ActorRef, ActorSystem, Behavior, PostStop, Scheduler, Signal, SpawnProtocol} +import org.apache.pekko.cluster.ddata.{ORSet, ORSetKey, SelfUniqueAddress} +import org.apache.pekko.actor.typed.scaladsl.AskPattern.* +import org.apache.pekko.cluster.ClusterEvent.{MemberEvent, ReachabilityEvent} +import org.apache.pekko.cluster.typed.{Cluster, Leave} +import org.apache.pekko.event.Logging.{DebugLevel, ErrorLevel, InfoLevel, LogLevel, WarningLevel} import vbds.server.routes.{AccessRoute, AdminRoute, TransferRoute} -import akka.http.scaladsl.server.Directives.* +import org.apache.pekko.http.scaladsl.server.Directives.* import csw.logging.client.scaladsl.GenericLoggerFactory import vbds.server.routes.AccessRoute.WebsocketResponseActorMsg @@ -126,11 +127,11 @@ private[server] object SharedDataActor { implicit actorSystem: ActorSystem[SpawnProtocol.Command] ): Behavior[SharedDataActorMessages] = { Behaviors.setup { ctx => - val memberEventAdapter: ActorRef[MemberEvent] = ctx.messageAdapter(MemberChange) - Cluster(ctx.system).subscriptions ! akka.cluster.typed.Subscribe(memberEventAdapter, classOf[MemberEvent]) + val memberEventAdapter: ActorRef[MemberEvent] = ctx.messageAdapter(MemberChange.apply) + Cluster(ctx.system).subscriptions ! org.apache.pekko.cluster.typed.Subscribe(memberEventAdapter, classOf[MemberEvent]) - val reachabilityAdapter = ctx.messageAdapter(ReachabilityChange) - Cluster(ctx.system).subscriptions ! akka.cluster.typed.Subscribe(reachabilityAdapter, classOf[ReachabilityEvent]) + val reachabilityAdapter = ctx.messageAdapter(ReachabilityChange.apply) + Cluster(ctx.system).subscriptions ! org.apache.pekko.cluster.typed.Subscribe(reachabilityAdapter, classOf[ReachabilityEvent]) DistributedData.withReplicatorMessageAdapter[SharedDataActorMessages, ORSet[?]] { replicatorAdapter => // Subscribe to changes of the given keys @@ -162,8 +163,8 @@ private[server] class SharedDataActor( // The local IP address, set at runtime via a message from the code that starts the HTTP server. // This is used to determine which HTTP server has the websocket connection to a client. - var localAddress: InetSocketAddress = _ - var binding: Http.ServerBinding = _ + var localAddress: InetSocketAddress = uninitialized + var binding: Http.ServerBinding = uninitialized // Cache of shared subscription info var subscriptions = Set[AccessInfo]() @@ -358,7 +359,7 @@ private[server] class SharedDataActor( ): Future[Done] = { // Split subscribers into local and remote - implicit val timeout = Timeout(20.seconds) + implicit val timeout: Timeout = Timeout(20.seconds) val (localSet, remoteSet) = getSubscribers(subscriberSet, dist) val remoteHostSet = remoteSet.map(a => ServerInfo(a.host, a.port)) if (dist) checkRemoteConnections(remoteHostSet) @@ -371,7 +372,7 @@ private[server] class SharedDataActor( // Wait for the client to acknowledge the message def waitForAck(a: AccessInfo): Future[Any] = { - localSubscribers(a).wsResponseActor.ask(AccessRoute.Get)(timeout, scheduler) + localSubscribers(a).wsResponseActor.ask(AccessRoute.Get.apply)(timeout, scheduler) } // Send data for a remote subscriber as HTTP POST to the server hosting its websocket @@ -433,7 +434,7 @@ private[server] class SharedDataActor( * @param dist true if published data should be distributed to remote subscribers */ private def getSubscribers(subscriberSet: Set[AccessInfo], dist: Boolean): (Set[AccessInfo], Set[AccessInfo]) = { - val (localSet, remoteSet) = subscriberSet.partition(localSubscribers.contains _) + val (localSet, remoteSet) = subscriberSet.partition(localSubscribers.contains) if (dist) (localSet, remoteSet) else (localSet, Set.empty[AccessInfo]) } diff --git a/vbds-server/src/main/scala/vbds/server/actors/TransferApi.scala b/vbds-server/src/main/scala/vbds/server/actors/TransferApi.scala index 46e78bd..7ad66d9 100644 --- a/vbds-server/src/main/scala/vbds/server/actors/TransferApi.scala +++ b/vbds-server/src/main/scala/vbds/server/actors/TransferApi.scala @@ -1,11 +1,11 @@ package vbds.server.actors -import akka.Done -import akka.actor.typed.{ActorRef, Scheduler} -import akka.stream.scaladsl.Source -import akka.util.{ByteString, Timeout} +import org.apache.pekko.Done +import org.apache.pekko.actor.typed.{ActorRef, Scheduler} +import org.apache.pekko.stream.scaladsl.Source +import org.apache.pekko.util.{ByteString, Timeout} import vbds.server.actors.SharedDataActor.{Publish, SharedDataActorMessages} -import akka.actor.typed.scaladsl.AskPattern._ +import org.apache.pekko.actor.typed.scaladsl.AskPattern._ import scala.concurrent.Future import scala.concurrent.duration._ diff --git a/vbds-server/src/main/scala/vbds/server/app/VbdsServer.scala b/vbds-server/src/main/scala/vbds/server/app/VbdsServer.scala index a0da1ab..92e34d2 100644 --- a/vbds-server/src/main/scala/vbds/server/app/VbdsServer.scala +++ b/vbds-server/src/main/scala/vbds/server/app/VbdsServer.scala @@ -1,12 +1,12 @@ package vbds.server.app -import akka.{Done, actor} -import akka.actor.{AddressFromURIString, CoordinatedShutdown} -import akka.actor.typed.scaladsl.AskPattern.Askable -import akka.actor.typed.scaladsl.adapter.TypedActorSystemOps -import akka.actor.typed.{SpawnProtocol, _} -import akka.cluster.typed.{Cluster, JoinSeedNodes} -import akka.util.Timeout +import org.apache.pekko.{Done, actor} +import org.apache.pekko.actor.{AddressFromURIString, CoordinatedShutdown} +import org.apache.pekko.actor.typed.scaladsl.AskPattern.Askable +import org.apache.pekko.actor.typed.scaladsl.adapter.TypedActorSystemOps +import org.apache.pekko.actor.typed.{SpawnProtocol, _} +import org.apache.pekko.cluster.typed.{Cluster, JoinSeedNodes} +import org.apache.pekko.util.Timeout import com.typesafe.config.ConfigFactory import csw.location.api.models.Connection.HttpConnection import csw.location.api.models.{ComponentId, ComponentType, HttpRegistration, Metadata, NetworkType} @@ -34,7 +34,7 @@ object VbdsServer { clusterSeeds .split(",") .toList - .map(s => s"akka://$clusterName@$s") + .map(s => s"pekko://$clusterName@$s") } /** @@ -42,7 +42,7 @@ object VbdsServer { * * @param host Akka and HTTP server bind host * @param httpPort HTTP server port (must be > 0) - * @param akkaPort akka ActorSystem port (must be > 0) + * @param pekkoPort akka ActorSystem port (must be > 0) * @param name the name of this server (for LocationService and ActorSystem) * @param clusterSeeds list of cluster seeds in the form host:port,host:port,... (Required even for seed node) * @return the root actor system @@ -50,17 +50,17 @@ object VbdsServer { def start( host: String, httpPort: Int, - akkaPort: Int, + pekkoPort: Int, name: String, clusterSeeds: String ): ActorSystem[SpawnProtocol.Command] = { checkPort(httpPort) - checkPort(akkaPort) + checkPort(pekkoPort) // Generate the akka config for the akka and http ports as well as the cluster seed nodes val config = ConfigFactory.parseString(s""" - akka.remote.artery.canonical.hostname=$host - akka.remote.artery.canonical.port=$akkaPort + pekko.remote.artery.canonical.hostname=$host + pekko.remote.artery.canonical.port=$pekkoPort """).withFallback(ConfigFactory.load()) implicit val system = ActorSystem(SpawnProtocol(), clusterName, config) @@ -125,7 +125,7 @@ object VbdsServer { .sequence( List( registrationResult.unregister(), - actorRef.ask(StopSharedDataActor)(Timeout(3.seconds), system.scheduler) + actorRef.ask(StopSharedDataActor.apply)(Timeout(3.seconds), system.scheduler) ) ) .map(_ => Done) diff --git a/vbds-server/src/main/scala/vbds/server/app/VbdsServerApp.scala b/vbds-server/src/main/scala/vbds/server/app/VbdsServerApp.scala index 07182a4..9eddaa8 100644 --- a/vbds-server/src/main/scala/vbds/server/app/VbdsServerApp.scala +++ b/vbds-server/src/main/scala/vbds/server/app/VbdsServerApp.scala @@ -6,13 +6,13 @@ import csw.network.utils.{Networks, SocketUtils} * VIZ Bulk Data System HTTP server and Akka cluster. * This is the command line app used to start the server. */ -object VbdsServerApp extends App { +object VbdsServerApp { // Command line options private case class Options( name: String = "vbds", httpPort: Int = 0, - akkaPort: Int = 0, + pekkoPort: Int = 0, clusterSeeds: String = "" ) @@ -28,9 +28,9 @@ object VbdsServerApp extends App { c.copy(httpPort = x) } text "The HTTP server port number (default: 0 for random port)" - opt[Int]("akka-port") valueName "" action { (x, c) => - c.copy(akkaPort = x) - } text "The Akka system port number (default: 0 for random port)" + opt[Int]("pekko-port") valueName "" action { (x, c) => + c.copy(pekkoPort = x) + } text "The Pekko system port number (default: 0 for random port)" opt[String]('s', "seeds") valueName ":,:,..." action { (x, c) => c.copy(clusterSeeds = x) @@ -41,26 +41,28 @@ object VbdsServerApp extends App { } // Parse the command line options - parser.parse(args, Options()) match { - case Some(options) => - try { - run(options) - } catch { - case e: Throwable => - e.printStackTrace() - System.exit(1) - } - case None => System.exit(1) + def main(args: Array[String]): Unit = { + parser.parse(args, Options()) match { + case Some(options) => + try { + run(options) + } catch { + case e: Throwable => + e.printStackTrace() + System.exit(1) + } + case None => System.exit(1) + } } // Run the application private def run(options: Options): Unit = { val host = Networks.publicInterface(None).hostname - val akkaPort = if (options.akkaPort != 0) options.akkaPort else SocketUtils.getFreePort + val pekkoPort = if (options.pekkoPort != 0) options.pekkoPort else SocketUtils.getFreePort val httpPort = if (options.httpPort != 0) options.httpPort else SocketUtils.getFreePort val clusterSeeds = if (options.clusterSeeds.nonEmpty) - options.clusterSeeds else s"${host}:$akkaPort" + options.clusterSeeds else s"${host}:$pekkoPort" - VbdsServer.start(host, httpPort, akkaPort, options.name, clusterSeeds) + VbdsServer.start(host, httpPort, pekkoPort, options.name, clusterSeeds) } } diff --git a/vbds-server/src/main/scala/vbds/server/marshalling/BFormat.scala b/vbds-server/src/main/scala/vbds/server/marshalling/BFormat.scala index 21035cf..9e5b437 100644 --- a/vbds-server/src/main/scala/vbds/server/marshalling/BFormat.scala +++ b/vbds-server/src/main/scala/vbds/server/marshalling/BFormat.scala @@ -1,6 +1,6 @@ package vbds.server.marshalling -import akka.util.ByteString +import org.apache.pekko.util.ByteString import boopickle.Default._ trait BFormat[T] { @@ -17,8 +17,8 @@ object BFormat { def write(o: T) = toBinary(o) } - implicit val stringFormat = make[String](_.utf8String, ByteString.apply) - implicit val byeStringFormat = make[ByteString](identity, identity) + implicit val stringFormat: BFormat[String] = make[String](_.utf8String, ByteString.apply) + implicit val byeStringFormat: BFormat[ByteString] = make[ByteString](identity, identity) implicit def objectFormat[T: Pickler]: BFormat[T] = BFormat.make[T]( x => Unpickle[T].fromBytes(x.toByteBuffer), diff --git a/vbds-server/src/main/scala/vbds/server/marshalling/BinaryMarshallers.scala b/vbds-server/src/main/scala/vbds/server/marshalling/BinaryMarshallers.scala index 29d099e..3c51d4d 100644 --- a/vbds-server/src/main/scala/vbds/server/marshalling/BinaryMarshallers.scala +++ b/vbds-server/src/main/scala/vbds/server/marshalling/BinaryMarshallers.scala @@ -1,15 +1,15 @@ package vbds.server.marshalling -import akka.http.scaladsl.marshalling.Marshaller -import akka.http.scaladsl.model.HttpEntity.Chunked -import akka.http.scaladsl.model.{ContentTypes, HttpEntity} -import akka.http.scaladsl.unmarshalling.Unmarshaller -import akka.stream.scaladsl.Source +import org.apache.pekko.http.scaladsl.marshalling.Marshaller +import org.apache.pekko.http.scaladsl.model.HttpEntity.Chunked +import org.apache.pekko.http.scaladsl.model.{ContentTypes, HttpEntity} +import org.apache.pekko.http.scaladsl.unmarshalling.Unmarshaller +import org.apache.pekko.stream.scaladsl.Source trait BinaryMarshallers { implicit def byteStringMarshaller[T: BFormat] : Marshaller[Source[T, Any], Chunked] = Marshaller.opaque { - source: Source[T, Any] => + (source: Source[T, Any]) => val byteStrings = source.map(BFormat[T].write) HttpEntity.Chunked.fromData(ContentTypes.`application/octet-stream`, byteStrings) @@ -17,7 +17,7 @@ trait BinaryMarshallers { implicit def byteStringUnmarshaller[T: BFormat] : Unmarshaller[HttpEntity, Source[T, Any]] = Unmarshaller.strict { - entity: HttpEntity => + (entity: HttpEntity) => entity.dataBytes.map(BFormat[T].read) } } diff --git a/vbds-server/src/main/scala/vbds/server/marshalling/VbdsSerializer.scala b/vbds-server/src/main/scala/vbds/server/marshalling/VbdsSerializer.scala index 2c219af..92b9b4c 100644 --- a/vbds-server/src/main/scala/vbds/server/marshalling/VbdsSerializer.scala +++ b/vbds-server/src/main/scala/vbds/server/marshalling/VbdsSerializer.scala @@ -1,11 +1,11 @@ package vbds.server.marshalling -import csw.commons.CborAkkaSerializer +import csw.commons.CborPekkoSerializer import vbds.server.models.{AccessInfo, ServerInfo, StreamInfo} import io.bullet.borer.Codec import io.bullet.borer.derivation.CompactMapBasedCodecs.deriveCodec -class VbdsSerializer extends CborAkkaSerializer[VbdsSerializable] { +class VbdsSerializer extends CborPekkoSerializer[VbdsSerializable] { override def identifier: Int = 19999 implicit lazy val accessInfoCodec: Codec[AccessInfo] = deriveCodec implicit lazy val serverInfoCodec: Codec[ServerInfo] = deriveCodec diff --git a/vbds-server/src/main/scala/vbds/server/models/JsonSupport.scala b/vbds-server/src/main/scala/vbds/server/models/JsonSupport.scala index a29aecc..98c1b3e 100644 --- a/vbds-server/src/main/scala/vbds/server/models/JsonSupport.scala +++ b/vbds-server/src/main/scala/vbds/server/models/JsonSupport.scala @@ -1,12 +1,23 @@ package vbds.server.models -import akka.http.scaladsl.marshallers.sprayjson.SprayJsonSupport -import spray.json.DefaultJsonProtocol +import org.apache.pekko.http.scaladsl.marshallers.sprayjson.SprayJsonSupport +import spray.json.* /** * Defines JSON I/O for model objects */ trait JsonSupport extends SprayJsonSupport with DefaultJsonProtocol { - implicit val streamInfoFormat = jsonFormat2(StreamInfo) - implicit val accessInfoFormat = jsonFormat4(AccessInfo) -} + implicit val streamInfoFormat: RootJsonFormat[StreamInfo] = jsonFormat2(StreamInfo.apply) + implicit val accessInfoFormat: RootJsonFormat[AccessInfo] = jsonFormat4(AccessInfo.apply) + + // spray-json's CollectionFormats.setFormat targets scala.collection.Set, which is + // incomparable to immutable.Iterable (the source of the Scala-3 ambiguity). A format + // typed to immutable Set is strictly the most specific candidate, so it wins resolution. + implicit def immutableSetFormat[T: JsonFormat]: RootJsonFormat[Set[T]] = new RootJsonFormat[Set[T]] { + def write(set: Set[T]): JsValue = JsArray(set.map(_.toJson).toVector) + def read(value: JsValue): Set[T] = value match { + case JsArray(elements) => elements.map(_.convertTo[T]).toSet + case x => deserializationError(s"Expected a Set as JsArray, but got $x") + } + } +} \ No newline at end of file diff --git a/vbds-server/src/main/scala/vbds/server/routes/AccessRoute.scala b/vbds-server/src/main/scala/vbds/server/routes/AccessRoute.scala index a4c8efe..d678ac4 100644 --- a/vbds-server/src/main/scala/vbds/server/routes/AccessRoute.scala +++ b/vbds-server/src/main/scala/vbds/server/routes/AccessRoute.scala @@ -2,16 +2,16 @@ package vbds.server.routes import java.util.UUID -import akka.actor.typed.{ActorRef, ActorSystem, Behavior, SpawnProtocol} -import akka.actor.typed.scaladsl.Behaviors -import akka.http.scaladsl.model.StatusCodes -import akka.http.scaladsl.model.ws.{BinaryMessage, Message} -import akka.http.scaladsl.server.Directives -import akka.util.ByteString +import org.apache.pekko.actor.typed.{ActorRef, ActorSystem, Behavior, SpawnProtocol} +import org.apache.pekko.actor.typed.scaladsl.Behaviors +import org.apache.pekko.http.scaladsl.model.StatusCodes +import org.apache.pekko.http.scaladsl.model.ws.{BinaryMessage, Message} +import org.apache.pekko.http.scaladsl.server.Directives +import org.apache.pekko.util.ByteString import vbds.server.actors.{AccessApi, AdminApi} import vbds.server.models.JsonSupport import AccessRoute._ -import akka.stream.scaladsl.{Flow, MergeHub, Sink} +import org.apache.pekko.stream.scaladsl.{Flow, MergeHub, Sink} import vbds.server.actors.AkkaTypedExtension.UserActorFactory // Actor to handle ACK responses from websocket clients @@ -23,13 +23,13 @@ object AccessRoute { final case class Get(replyTo: ActorRef[Ack.type]) extends WebsocketResponseActorMsg // Says there was a response from the ws client - final case object Put extends WebsocketResponseActorMsg + case object Put extends WebsocketResponseActorMsg // Stops the actor - final case object Stop extends WebsocketResponseActorMsg + case object Stop extends WebsocketResponseActorMsg // Reponse to Get message - final case object Ack + case object Ack // Actor that handles responses from the websocket private def websocketResponseBehavior( diff --git a/vbds-server/src/main/scala/vbds/server/routes/AdminRoute.scala b/vbds-server/src/main/scala/vbds/server/routes/AdminRoute.scala index 26afbb7..41c3161 100644 --- a/vbds-server/src/main/scala/vbds/server/routes/AdminRoute.scala +++ b/vbds-server/src/main/scala/vbds/server/routes/AdminRoute.scala @@ -1,7 +1,7 @@ package vbds.server.routes -import akka.http.scaladsl.model.StatusCodes._ -import akka.http.scaladsl.server.Directives +import org.apache.pekko.http.scaladsl.model.StatusCodes._ +import org.apache.pekko.http.scaladsl.server.Directives import vbds.server.actors.AdminApi import vbds.server.models.JsonSupport diff --git a/vbds-server/src/main/scala/vbds/server/routes/Cors.scala b/vbds-server/src/main/scala/vbds/server/routes/Cors.scala index 7e2cd89..d3a6072 100644 --- a/vbds-server/src/main/scala/vbds/server/routes/Cors.scala +++ b/vbds-server/src/main/scala/vbds/server/routes/Cors.scala @@ -1,10 +1,10 @@ package vbds.server.routes -import akka.http.scaladsl.model.HttpMethods._ -import akka.http.scaladsl.model.{HttpHeader, HttpResponse} -import akka.http.scaladsl.model.headers._ -import akka.http.scaladsl.server.Directives._ -import akka.http.scaladsl.server.{Directive0, MethodRejection, RejectionHandler} +import org.apache.pekko.http.scaladsl.model.HttpMethods._ +import org.apache.pekko.http.scaladsl.model.{HttpHeader, HttpResponse} +import org.apache.pekko.http.scaladsl.model.headers._ +import org.apache.pekko.http.scaladsl.server.Directives._ +import org.apache.pekko.http.scaladsl.server.{Directive0, MethodRejection, RejectionHandler} object Cors { val corsAllowOrigins: List[String] = List("*") diff --git a/vbds-server/src/main/scala/vbds/server/routes/LoggingSupport.scala b/vbds-server/src/main/scala/vbds/server/routes/LoggingSupport.scala index 0288f13..69882af 100644 --- a/vbds-server/src/main/scala/vbds/server/routes/LoggingSupport.scala +++ b/vbds-server/src/main/scala/vbds/server/routes/LoggingSupport.scala @@ -1,7 +1,7 @@ package vbds.server.routes -import akka.actor.typed.{ActorSystem, SpawnProtocol} -import akka.event.{LogSource, Logging} +import org.apache.pekko.actor.typed.{ActorSystem, SpawnProtocol} +import org.apache.pekko.event.{LogSource, Logging} /** * Logging support for routes @@ -11,7 +11,7 @@ trait LoggingSupport { implicit val logSource: LogSource[AnyRef] = new LogSource[AnyRef] { def genString(o: AnyRef): String = o.getClass.getName - override def getClazz(o: AnyRef): Class[_] = o.getClass + override def getClazz(o: AnyRef): Class[?] = o.getClass } val log = Logging(actorSystem.classicSystem, this) diff --git a/vbds-server/src/main/scala/vbds/server/routes/TransferRoute.scala b/vbds-server/src/main/scala/vbds/server/routes/TransferRoute.scala index dbd96a7..ec1796e 100644 --- a/vbds-server/src/main/scala/vbds/server/routes/TransferRoute.scala +++ b/vbds-server/src/main/scala/vbds/server/routes/TransferRoute.scala @@ -1,13 +1,13 @@ package vbds.server.routes -import akka.actor.typed.{ActorSystem, SpawnProtocol} -import akka.http.scaladsl.model.StatusCodes._ -import akka.http.scaladsl.server.Directives +import org.apache.pekko.actor.typed.{ActorSystem, SpawnProtocol} +import org.apache.pekko.http.scaladsl.model.StatusCodes._ +import org.apache.pekko.http.scaladsl.server.Directives import vbds.server.actors.{AdminApi, TransferApi} import vbds.server.models.JsonSupport import vbds.server.marshalling.BinaryMarshallers -import akka.stream.scaladsl.Source -import akka.util.ByteString +import org.apache.pekko.stream.scaladsl.Source +import org.apache.pekko.util.ByteString /** * Provides the HTTP route for the VBDS Transfer Service. diff --git a/vbds-server/src/multi-jvm/scala/vbds/server/VbdsServerTest.scala b/vbds-server/src/multi-jvm/scala/vbds/server/VbdsServerTest.scala index edb8f6a..dff6d01 100644 --- a/vbds-server/src/multi-jvm/scala/vbds/server/VbdsServerTest.scala +++ b/vbds-server/src/multi-jvm/scala/vbds/server/VbdsServerTest.scala @@ -1,17 +1,17 @@ package vbds.server import java.io.{BufferedOutputStream, File, FileOutputStream} -import akka.actor.{Actor, ActorLogging, PoisonPill, Props} +import org.apache.pekko.actor.{Actor, ActorLogging, PoisonPill, Props} import vbds.client.VbdsClient import vbds.server.app.VbdsServer import scala.concurrent.duration.{Duration, DurationLong, FiniteDuration} import scala.concurrent.{Await, Future, Promise} -import akka.remote.testkit.{MultiNodeConfig, MultiNodeSpec} -import akka.testkit.ImplicitSender -import akka.event.LoggingAdapter -import akka.http.scaladsl.model.StatusCodes -import akka.stream.scaladsl.{Sink, Source} +import org.apache.pekko.remote.testkit.{MultiNodeConfig, MultiNodeSpec} +import org.apache.pekko.testkit.ImplicitSender +import org.apache.pekko.event.LoggingAdapter +import org.apache.pekko.http.scaladsl.model.StatusCodes +import org.apache.pekko.stream.scaladsl.{Sink, Source} import com.typesafe.config.{Config, ConfigFactory} import csw.network.utils.SocketUtils import org.apache.commons.io.FileUtils @@ -37,10 +37,10 @@ object VbdsServerTestConfig extends MultiNodeConfig { val configStr = """ - | akka.loglevel = INFO - | akka.log-dead-letters-during-shutdown = off - | akka.testconductor.barrier-timeout = 30m - | akka.cluster.jmx.multi-mbeans-in-same-jvm = on + | pekko.loglevel = INFO + | pekko.log-dead-letters-during-shutdown = off + | pekko.testconductor.barrier-timeout = 30m + | pekko.cluster.jmx.multi-mbeans-in-same-jvm = on | remote { | artery { | enabled = on @@ -53,8 +53,8 @@ object VbdsServerTestConfig extends MultiNodeConfig { commonConfig(ConfigFactory.parseString(configStr).withFallback(ConfigFactory.load())) -// def makeSystem(config: Config): akka.actor.ActorSystem = ActorSystem(SpawnProtocol(), VbdsServer.clusterName, config).classicSystem - def makeSystem(config: Config) = akka.actor.ActorSystem(VbdsServer.clusterName, config) +// def makeSystem(config: Config): org.apache.pekko.actor.ActorSystem = ActorSystem(SpawnProtocol(), VbdsServer.clusterName, config).classicSystem + def makeSystem(config: Config) = org.apache.pekko.actor.ActorSystem(VbdsServer.clusterName, config) } // One for each role @@ -317,4 +317,4 @@ class VbdsServerTest enterBarrier("finished") log.debug(s"enterBarrier finished") } -} +} \ No newline at end of file diff --git a/vbds-server/src/test/scala/vbds/server/STMultiNodeSpec.scala b/vbds-server/src/test/scala/vbds/server/STMultiNodeSpec.scala index 1324acf..49c0ddf 100644 --- a/vbds-server/src/test/scala/vbds/server/STMultiNodeSpec.scala +++ b/vbds-server/src/test/scala/vbds/server/STMultiNodeSpec.scala @@ -1,6 +1,6 @@ package vbds.server -import akka.remote.testkit.{MultiNodeSpec, MultiNodeSpecCallbacks} +import org.apache.pekko.remote.testkit.{MultiNodeSpec, MultiNodeSpecCallbacks} import scala.language.implicitConversions import org.scalatest.BeforeAndAfterAll @@ -20,4 +20,4 @@ trait STMultiNodeSpec extends MultiNodeSpecCallbacks with AnyWordSpecLike with M // Might not be needed anymore if we find a nice way to tag all logging from a node override implicit def convertToWordSpecStringWrapper(s: String): WordSpecStringWrapper = new WordSpecStringWrapper(s"$s (on node '${self.myself.name}', $getClass)") -} +} \ No newline at end of file diff --git a/vbds-server/src/test/scala/vbds/server/routes/LocalAdminApi.scala b/vbds-server/src/test/scala/vbds/server/routes/LocalAdminApi.scala index 056458b..2d2aff9 100644 --- a/vbds-server/src/test/scala/vbds/server/routes/LocalAdminApi.scala +++ b/vbds-server/src/test/scala/vbds/server/routes/LocalAdminApi.scala @@ -1,6 +1,6 @@ package vbds.server.routes -import akka.actor.ActorSystem +import org.apache.pekko.actor.ActorSystem import vbds.server.actors.AdminApi import vbds.server.models.StreamInfo @@ -30,4 +30,4 @@ class LocalAdminApi(system: ActorSystem) extends AdminApi { streams = streams - result Future.successful(result) } -} +} \ No newline at end of file diff --git a/vbds-server/src/test/scala/vbds/server/routes/StreamAdminSpec.scala b/vbds-server/src/test/scala/vbds/server/routes/StreamAdminSpec.scala index de15da8..6b7b42b 100644 --- a/vbds-server/src/test/scala/vbds/server/routes/StreamAdminSpec.scala +++ b/vbds-server/src/test/scala/vbds/server/routes/StreamAdminSpec.scala @@ -1,9 +1,9 @@ package vbds.server.routes -import akka.http.scaladsl.model.StatusCodes +import org.apache.pekko.http.scaladsl.model.StatusCodes import org.scalatest.wordspec.AnyWordSpec import org.scalatest.matchers.should.Matchers -import akka.http.scaladsl.testkit.ScalatestRouteTest +import org.apache.pekko.http.scaladsl.testkit.ScalatestRouteTest import vbds.server.models.{JsonSupport, StreamInfo} class StreamAdminSpec extends AnyWordSpec with Matchers with ScalatestRouteTest with JsonSupport { @@ -57,4 +57,4 @@ class StreamAdminSpec extends AnyWordSpec with Matchers with ScalatestRouteTest } } } -} +} \ No newline at end of file