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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -16,7 +16,6 @@ package org.apache.pekko.discovery.awsapi.ecs
import java.net.InetAddress
import java.util.concurrent.TimeoutException

import scala.collection.immutable.Seq
import scala.concurrent.duration._
import scala.concurrent.{ ExecutionContext, Future }
import scala.jdk.CollectionConverters._
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,7 +16,6 @@ package org.apache.pekko.discovery.awsapi.ecs
import java.net.InetAddress
import java.util.concurrent.TimeoutException

import scala.collection.immutable.Seq
import scala.concurrent.duration._
import scala.concurrent.{ ExecutionContext, Future }
import scala.jdk.CollectionConverters._
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -30,7 +30,6 @@ import pekko.pattern.after
import java.net.InetAddress
import java.util.concurrent.TimeoutException
import scala.annotation.tailrec
import scala.collection.immutable.Seq
import scala.concurrent.duration.FiniteDuration
import scala.concurrent.{ ExecutionContext, Future }
import scala.jdk.CollectionConverters._
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -31,7 +31,6 @@ import com.amazonaws.services.ecs.model.{ DescribeTasksRequest, DesiredStatus, L
import com.amazonaws.services.ecs.{ AmazonECS, AmazonECSClientBuilder }

import scala.annotation.tailrec
import scala.collection.immutable.Seq
import scala.concurrent.{ ExecutionContext, Future }
import scala.concurrent.duration._
import scala.jdk.CollectionConverters._
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -34,7 +34,6 @@ import java.security.cert.CertificateFactory
import java.util
import java.util.concurrent.TimeoutException
import javax.net.ssl.{ SSLContext, TrustManagerFactory }
import scala.collection.immutable.Seq
import scala.concurrent.duration.FiniteDuration
import scala.concurrent.{ ExecutionContext, Future, Promise }
import scala.jdk.CollectionConverters._
Expand Down Expand Up @@ -107,7 +106,7 @@ class ConsulServiceDiscovery(system: ActorSystem) extends ServiceDiscovery {
Future(extractResolvedTargetFromCatalogService(catalogService))(blockingEc)
}
} yield resolvedTargets
consulResult.map(targets => Resolved(name, scala.collection.immutable.Seq(targets: _*)))
consulResult.map(targets => Resolved(name, targets))
}

private def extractResolvedTargetFromCatalogService(catalogService: CatalogService) = {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,8 +18,6 @@ import java.nio.charset.StandardCharsets
import java.util.concurrent.TimeoutException
import java.nio.file.{ Files, Paths }

import scala.collection.immutable
import scala.collection.immutable.Seq
import scala.concurrent.ExecutionContext
import scala.concurrent.Future
import scala.concurrent.Promise
Expand Down Expand Up @@ -59,7 +57,7 @@ object KubernetesApiServiceDiscovery {
podNamespace: String,
podDomain: String,
rawIp: Boolean,
containerName: Option[String]): immutable.Seq[ResolvedTarget] =
containerName: Option[String]): Seq[ResolvedTarget] =
for {
item <- podList.items
if item.metadata.flatMap(_.deletionTimestamp).isEmpty
Expand Down Expand Up @@ -240,7 +238,7 @@ class KubernetesApiServiceDiscovery(settings: Settings)(
val query = Uri.Query("labelSelector" -> labelSelector)
val uri = Uri.from(scheme = "https", host = host, port = port).withPath(path).withQuery(query)

val authHeaders = immutable.Seq(Authorization(OAuth2BearerToken(token)))
val authHeaders = Seq(Authorization(OAuth2BearerToken(token)))
val acceptEncodingHeader = HttpEncodings.getForKey(settings.httpRequestAcceptEncoding)
.map(httpEncoding => AcceptEncoding.create(httpEncoding))
HttpRequest(uri = uri, headers = authHeaders ++ acceptEncodingHeader)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,6 @@

package org.apache.pekko.discovery.kubernetes

import scala.collection.immutable
import org.apache.pekko.annotation.InternalApi

/**
Expand All @@ -24,15 +23,15 @@ import org.apache.pekko.annotation.InternalApi

final case class ContainerPort(name: Option[String], containerPort: Int)

final case class Container(name: String, ports: Option[immutable.Seq[ContainerPort]])
final case class Container(name: String, ports: Option[Seq[ContainerPort]])

final case class PodSpec(containers: immutable.Seq[Container])
final case class PodSpec(containers: Seq[Container])

final case class ContainerStatus(name: String, state: Map[String, Unit])

final case class PodStatus(
podIP: Option[String],
containerStatuses: Option[immutable.Seq[ContainerStatus]],
containerStatuses: Option[Seq[ContainerStatus]],
phase: Option[String])

final case class Pod(spec: Option[PodSpec], status: Option[PodStatus], metadata: Option[Metadata])
Expand All @@ -41,4 +40,4 @@ import org.apache.pekko.annotation.InternalApi
/**
* INTERNAL API
*/
@InternalApi private[kubernetes] final case class PodList(items: immutable.Seq[PodList.Pod])
@InternalApi private[kubernetes] final case class PodList(items: Seq[PodList.Pod])
Original file line number Diff line number Diff line change
Expand Up @@ -13,8 +13,6 @@

package org.apache.pekko.discovery.marathon

import scala.collection.immutable.Seq

object AppList {
case class App(container: Option[Container], portDefinitions: Option[Seq[PortDefinition]], tasks: Option[Seq[Task]])
case class Container(portMappings: Option[Seq[PortMapping]], docker: Option[Docker])
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,6 @@ import pekko.http.scaladsl._
import pekko.http.scaladsl.model._
import pekko.http.scaladsl.unmarshalling.Unmarshal

import scala.collection.immutable.Seq
import scala.concurrent.Future
import scala.concurrent.duration.FiniteDuration
import scala.util.Try
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,6 @@ package org.apache.pekko.coordination.lease.kubernetes

import java.util.concurrent.Executors

import scala.collection.immutable
import scala.concurrent.ExecutionContext
import scala.concurrent.Future

Expand Down Expand Up @@ -76,7 +75,7 @@ class LeaseContentionSpec extends TestKit(ActorSystem("LeaseContentionSpec",
val nrClients = 30
implicit val ec: ExecutionContext = ExecutionContext.fromExecutor(Executors.newFixedThreadPool(nrClients)) // too many = HTTP request queue of pool fills up
// could make this more contended with a countdown latch so they all start at the same time
val leases: immutable.Seq[(String, Boolean)] = Future.sequence((0 until nrClients).map(i => {
val leases: Seq[(String, Boolean)] = Future.sequence((0 until nrClients).map(i => {
val clientName = s"client$i"
val lease = underTest.getLease(lease1, KubernetesLease.configPath, clientName)
Future {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -31,7 +31,6 @@ import pekko.util.ByteString

import java.nio.file.{ Files, Paths }
import javax.net.ssl.SSLContext
import scala.collection.immutable
import scala.concurrent.{ ExecutionContext, Future, Promise }
import scala.util.control.NonFatal

Expand Down Expand Up @@ -69,7 +68,7 @@ import scala.util.control.NonFatal
_.getOrElse(""))(ExecutionContext.parasitic)
private def headers() = if (settings.secure) {
apiToken().map { token =>
immutable.Seq(Authorization(OAuth2BearerToken(token)))
Seq(Authorization(OAuth2BearerToken(token)))
}(ExecutionContext.parasitic)
} else
Future.successful(Nil)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -85,16 +85,16 @@ private[pekko] object BootstrapCoordinator {
lookup: Lookup,
fallbackPort: Int,
filterOnFallbackPort: Boolean,
contactPoints: immutable.Seq[ResolvedTarget]): immutable.Iterable[ResolvedTarget] = {
contactPoints: Seq[ResolvedTarget]): immutable.Iterable[ResolvedTarget] = {

// if the user has specified a port name in the search, don't do any filtering and assume it
// is handled in the service discovery mechanism
if (lookup.portName.isDefined || !filterOnFallbackPort) {
contactPoints
} else {
contactPoints.groupBy(_.host).flatMap {
case (_, immutable.Seq(singleResult)) =>
immutable.Seq(singleResult)
case (_, Seq(singleResult)) =>
Seq(singleResult)
case (_, multipleResults) =>
if (multipleResults.exists(_.port.isDefined)) {
multipleResults.filter(_.port.contains(fallbackPort))
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -20,22 +20,22 @@ import spray.json.{ DefaultJsonProtocol, RootJsonFormat }

import scala.collection.immutable

final case class ClusterUnreachableMember(node: String, observedBy: immutable.Seq[String])
final case class ClusterUnreachableMember(node: String, observedBy: Seq[String])
final case class ClusterMember(node: String, nodeUid: String, status: String, roles: Set[String])
object ClusterMember {
implicit val clusterMemberOrdering: Ordering[ClusterMember] = Ordering.by(_.node)
}
final case class ClusterMembers(
selfNode: String,
members: Set[ClusterMember],
unreachable: immutable.Seq[ClusterUnreachableMember],
unreachable: Seq[ClusterUnreachableMember],
leader: Option[String],
oldest: Option[String],
oldestPerRole: Map[String, String])
final case class ClusterHttpManagementMessage(message: String)
final case class ShardEntityTypeKeys(entityTypeKeys: immutable.Set[String])
final case class ShardRegionInfo(shardId: String, numEntities: Int)
final case class ShardDetails(regions: immutable.Seq[ShardRegionInfo])
final case class ShardDetails(regions: Seq[ShardRegionInfo])

/** INTERNAL API */
@InternalApi private[pekko] sealed trait ClusterHttpManagementMemberOperation
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -41,7 +41,7 @@ import org.scalatest.matchers.should.Matchers
import org.scalatest.time.{ Millis, Seconds, Span }
import org.scalatest.wordspec.AnyWordSpecLike

import scala.collection.immutable._
import scala.collection.immutable.SortedSet
import scala.concurrent.Promise

class ClusterHttpManagementRoutesSpec
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,6 @@ package org.apache.pekko.management

import com.typesafe.config.Config

import scala.collection.immutable
import scala.concurrent.duration.FiniteDuration
import scala.jdk.CollectionConverters._
import scala.jdk.DurationConverters._
Expand Down Expand Up @@ -119,18 +118,18 @@ object HealthCheckSettings {
* @param checkTimeout how long to wait for all health checks to complete
*/
final class HealthCheckSettings(
val startupChecks: immutable.Seq[NamedHealthCheck],
val readinessChecks: immutable.Seq[NamedHealthCheck],
val livenessChecks: immutable.Seq[NamedHealthCheck],
val startupChecks: Seq[NamedHealthCheck],
val readinessChecks: Seq[NamedHealthCheck],
val livenessChecks: Seq[NamedHealthCheck],
val startupPath: String,
val readinessPath: String,
val livenessPath: String,
val checkTimeout: FiniteDuration) {

@deprecated("Use constructor that takes `startupChecks` and `startupPath` parameters instead", "1.1.0")
def this(
readinessChecks: immutable.Seq[NamedHealthCheck],
livenessChecks: immutable.Seq[NamedHealthCheck],
readinessChecks: Seq[NamedHealthCheck],
livenessChecks: Seq[NamedHealthCheck],
readinessPath: String,
livenessPath: String,
checkTimeout: FiniteDuration
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,7 +16,6 @@ package org.apache.pekko.management
import java.net.InetAddress
import java.util.Optional

import scala.collection.immutable
import scala.concurrent.duration.{ Duration, FiniteDuration }
import scala.jdk.CollectionConverters._
import scala.jdk.DurationConverters._
Expand Down Expand Up @@ -62,7 +61,7 @@ final class PekkoManagementSettings(val config: Config) {
val BasePath: Option[String] =
Option(cc.getString("base-path")).flatMap(it => if (it.trim == "") None else Some(it))

val RouteProviders: immutable.Seq[NamedRouteProvider] = {
val RouteProviders: Seq[NamedRouteProvider] = {
def validFQCN(value: Any) = {
value != null &&
value != "null" &&
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -28,7 +28,6 @@ import pekko.management.javadsl.{ ReadinessCheckSetup => JReadinessCheckSetup }
import pekko.management.javadsl.{ StartupCheckSetup => JStartupCheckSetup }
import pekko.management.scaladsl.{ HealthChecks, LivenessCheckSetup, ReadinessCheckSetup, StartupCheckSetup }

import scala.collection.immutable
import scala.concurrent.Future
import scala.jdk.CollectionConverters._
import scala.jdk.FutureConverters._
Expand Down Expand Up @@ -59,7 +58,7 @@ final private[pekko] class HealthChecksImpl(system: ExtendedActorSystem, setting
"Loading liveness checks [{}]",
settings.livenessChecks.map(a => a.name -> a.fullyQualifiedClassName).mkString(", "))

private val startupChecks: immutable.Seq[HealthCheck] = {
private val startupChecks: Seq[HealthCheck] = {
val fromScaladslSetup = system.settings.setup.get[StartupCheckSetup] match {
case None => Nil
case Some(setup) => setup.createHealthChecks(system)
Expand All @@ -72,7 +71,7 @@ final private[pekko] class HealthChecksImpl(system: ExtendedActorSystem, setting
fromConfig ++ fromScaladslSetup ++ fromJavadslSetup
}

private val readiness: immutable.Seq[HealthCheck] = {
private val readiness: Seq[HealthCheck] = {
val fromScaladslSetup = system.settings.setup.get[ReadinessCheckSetup] match {
case None => Nil
case Some(setup) => setup.createHealthChecks(system)
Expand All @@ -85,7 +84,7 @@ final private[pekko] class HealthChecksImpl(system: ExtendedActorSystem, setting
fromConfig ++ fromScaladslSetup ++ fromJavadslSetup
}

private val liveness: immutable.Seq[HealthCheck] = {
private val liveness: Seq[HealthCheck] = {
val fromScaladslSetup = system.settings.setup.get[LivenessCheckSetup] match {
case None => Nil
case Some(setup) => setup.createHealthChecks(system)
Expand All @@ -99,7 +98,7 @@ final private[pekko] class HealthChecksImpl(system: ExtendedActorSystem, setting
}

private def convertSuppliersToScala(
suppliers: JList[Supplier[CompletionStage[JBoolean]]]): immutable.Seq[HealthCheck] = {
suppliers: JList[Supplier[CompletionStage[JBoolean]]]): Seq[HealthCheck] = {
suppliers.asScala.toList.map(convertSupplierToScala)
}

Expand All @@ -111,7 +110,7 @@ final private[pekko] class HealthChecksImpl(system: ExtendedActorSystem, setting
system.dynamicAccess
.createInstanceFor[HealthCheck](
fqcn,
immutable.Seq((classOf[ActorSystem], system)))
Seq((classOf[ActorSystem], system)))
.recoverWith {
case _: NoSuchMethodException =>
system.dynamicAccess.createInstanceFor[HealthCheck](fqcn, Nil)
Expand All @@ -122,7 +121,7 @@ final private[pekko] class HealthChecksImpl(system: ExtendedActorSystem, setting
system.dynamicAccess
.createInstanceFor[Supplier[CompletionStage[JBoolean]]](
fqcn,
immutable.Seq((classOf[ActorSystem], system)))
Seq((classOf[ActorSystem], system)))
.recoverWith {
case _: NoSuchMethodException =>
system.dynamicAccess.createInstanceFor[Supplier[CompletionStage[JBoolean]]](fqcn, Nil)
Expand All @@ -131,7 +130,7 @@ final private[pekko] class HealthChecksImpl(system: ExtendedActorSystem, setting
}

private def load(
checks: immutable.Seq[NamedHealthCheck]): immutable.Seq[HealthCheck] = {
checks: Seq[NamedHealthCheck]): Seq[HealthCheck] = {
checks
.map(namedHealthCheck =>
tryLoadScalaHealthCheck(namedHealthCheck.fullyQualifiedClassName).recoverWith {
Expand Down Expand Up @@ -202,7 +201,7 @@ final private[pekko] class HealthChecksImpl(system: ExtendedActorSystem, setting
Future.fromTry(Try(check())).flatMap(identity)
}

private def check(checks: immutable.Seq[HealthCheck]): Future[Either[String, Unit]] = {
private def check(checks: Seq[HealthCheck]): Future[Either[String, Unit]] = {
val spawnedChecks: Seq[Future[Either[String, Unit]]] = checks.map { check =>
val checkName = check.getClass.getName
// Create a per-check timeout so each check gets its own timer,
Expand Down
Loading