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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion build.sbt
Original file line number Diff line number Diff line change
Expand Up @@ -215,6 +215,7 @@ lazy val tests = crossProject(JVMPlatform, JSPlatform, NativePlatform)
"eu.timepit" %%% "refined-cats" % refinedVersion,
"org.typelevel" %%% "otel4s-sdk-trace" % otel4sSdkVersion,
"org.typelevel" %%% "otel4s-sdk-exporter-trace" % otel4sSdkVersion,
"org.typelevel" %%% "otel4s-sdk-testkit" % otel4sSdkVersion,
),
testFrameworks += new TestFramework("munit.Framework"),
testOptions += {
Expand All @@ -232,7 +233,6 @@ lazy val tests = crossProject(JVMPlatform, JSPlatform, NativePlatform)
.jvmSettings(
Test / fork := true,
javaOptions += "-Dotel.service.name=SkunkTests",
libraryDependencies += "org.typelevel" %% "otel4s-sdk-testkit" % otel4sSdkVersion % Test,
)
.jsSettings(
scalaJSLinkerConfig ~= { _.withESFeatures(_.withESVersion(org.scalajs.linker.interface.ESVersion.ES2018)) },
Expand Down
11 changes: 6 additions & 5 deletions modules/core/shared/src/main/scala/Session.scala
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,7 @@ import skunk.net.SSLNegotiation
import skunk.net.protocol.Describe
import scala.concurrent.duration.Duration
import skunk.net.protocol.Parse
import skunk.telemetry.{ConnectionInfo, Telemetry, TelemetryConfig}
import skunk.telemetry.{ConnectionInfo, PoolTelemetry, Telemetry, TelemetryConfig}

/**
* Represents a live connection to a Postgres database. Operations provided here are safe to use
Expand Down Expand Up @@ -649,15 +649,16 @@ object Session {
def pooled(max: Int): Resource[F, Resource[F, Session[F]]] =
for {
telemetry <- Resource.eval(Telemetry.create(telemetryConfig, connectionInfo(database.getOrElse(""))))
pool <- pooledWithTelemetry(max)
} yield pool(telemetry)
pool <- pooledWithTelemetry(max)(telemetry)
} yield pool

private def pooledWithTelemetry(max: Int): Resource[F, Telemetry[F] => Resource[F, Session[F]]] = {
private def pooledWithTelemetry(max: Int)(implicit T: Telemetry[F]): Resource[F, Resource[F, Session[F]]] = {
implicit val poolTelemetry: PoolTelemetry[F] = T.pool
val logger: String => F[Unit] = s => Console[F].println(s"TLS: $s")
for {
dc <- Resource.eval(Describe.Cache.empty[F](commandCacheSize, queryCacheSize))
sslOp <- ssl.toSSLNegotiationOptions(if (debug) logger.some else none)
pool <- Pool.ofF({implicit T: Telemetry[F] => sessions(sslOp, dc)}, max, checkout = Recyclers.ensureHealthy[F], checkin = Recyclers.full)
pool <- Pool.of(sessions(sslOp, dc), max, checkout = Recyclers.ensureHealthy[F], checkin = Recyclers.full)
} yield pool
}

Expand Down
26 changes: 21 additions & 5 deletions modules/core/shared/src/main/scala/telemetry/ConnectionInfo.scala
Original file line number Diff line number Diff line change
Expand Up @@ -4,8 +4,24 @@

package skunk.telemetry

private[skunk] final case class ConnectionInfo(
database: String,
serverAddress: String,
serverPort: Option[Long]
)
sealed trait ConnectionInfo {

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Since Telemetry.create is public, ConnectionInfo must also be public.

def database: String
def serverAddress: String
def serverPort: Option[Long]
}

object ConnectionInfo {

def apply(
database: String,
serverAddress: String,
serverPort: Option[Long]
): ConnectionInfo =
Impl(database, serverAddress, serverPort)

private final case class Impl(
database: String,
serverAddress: String,
serverPort: Option[Long]
) extends ConnectionInfo
}
92 changes: 92 additions & 0 deletions modules/core/shared/src/main/scala/telemetry/PoolTelemetry.scala
Original file line number Diff line number Diff line change
@@ -0,0 +1,92 @@
// Copyright (c) 2018-2024 by Rob Norris and Contributors
// This software is licensed under the MIT License (MIT).
// For more information see LICENSE or https://opensource.org/licenses/MIT

package skunk.telemetry

import cats.arrow.FunctionK
import cats.syntax.functor._
import cats.{Monad, ~>}
import org.typelevel.otel4s.trace.{SpanKind, Tracer, TracerProvider}
import skunk.BuildInfo

sealed trait PoolTelemetry[F[_]] {

private[skunk] def span[A](name: String)(fa: F[A]): F[A]

}

object PoolTelemetry {

sealed trait Config {
def poolSpans: Config.PoolSpans

/** Enables or disables connection-pool spans. */
def withPoolSpans(poolSpans: Config.PoolSpans): Config
}

object Config {

/** Controls spans emitted by Skunk's connection pool. */
sealed trait PoolSpans
object PoolSpans {

/** Emit connection-pool operations as `INTERNAL` spans. */
case object Internal extends PoolSpans

/** Do not export connection-pool spans. */
case object Disabled extends PoolSpans
}

/** Recommended defaults. Pool spans are disabled.
*/
val default: Config = Config(PoolSpans.Disabled)

/** Creates a telemetry configuration with explicit settings. */
def apply(poolSpans: PoolSpans): Config =
Impl(poolSpans)

private final case class Impl(poolSpans: PoolSpans) extends Config {
def withPoolSpans(poolSpans: Config.PoolSpans): Config =
copy(poolSpans = poolSpans)

override def toString: String =
s"PoolTelemetry.Config($poolSpans)"
}

}

def apply[F[_]](implicit ev: PoolTelemetry[F]): PoolTelemetry[F] = ev

def create[F[_]: Monad: TracerProvider](config: Config): F[PoolTelemetry[F]] =
TracerProvider[F].tracer("org.typelevel.skunk").withVersion(BuildInfo.version).get.map { implicit tracer =>
new Impl(config)
}

def noop[F[_]]: PoolTelemetry[F] =
new Noop[F]

private final class Impl[F[_]: Tracer](config: Config) extends PoolTelemetry[F] {
private val spanF: String => F ~> F =
config.poolSpans match {
case Config.PoolSpans.Internal =>
label =>
FunctionK.liftFunction[F, F](
Tracer[F]
.spanBuilder(label)
.withSpanKind(SpanKind.Internal)
.build
.surround
)
case Config.PoolSpans.Disabled =>
Function.const(FunctionK.id[F])(_)
}

private[skunk] def span[A](name: String)(fa: F[A]): F[A] = spanF(name)(fa)
}

private final class Noop[F[_]] extends PoolTelemetry[F] {
def span[A](name: String)(fa: F[A]): F[A] = fa
}

}
59 changes: 34 additions & 25 deletions modules/core/shared/src/main/scala/telemetry/Telemetry.scala
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,7 @@ import cats.effect.{MonadCancelThrow, Resource}
import cats.syntax.flatMap._
import cats.syntax.functor._
import cats.syntax.semigroup._
import cats.~>
import cats.{Applicative, ~>}
import org.typelevel.otel4s.{Attribute, Attributes}
import org.typelevel.otel4s.metrics.{BucketBoundaries, Histogram, Meter, MeterProvider}
import org.typelevel.otel4s.semconv.attributes.{DbAttributes, ErrorAttributes, ServerAttributes}
Expand All @@ -27,8 +27,6 @@ sealed trait Telemetry[F[_]] {

private[skunk] def withConnection(connection: ConnectionInfo): Telemetry[F]

private[skunk] def poolSpan[A](name: String)(fa: F[A]): F[A]

private[skunk] def internalSpan[A](label: String)(fa: F[A]): F[A]

private[skunk] def databaseSpan[A](
Expand All @@ -42,6 +40,8 @@ sealed trait Telemetry[F[_]] {

private[skunk] def addProtocolAttributes(attributes: Attribute[_]*): F[Unit]

private[skunk] def pool: PoolTelemetry[F]

}

object Telemetry {
Expand All @@ -68,18 +68,23 @@ object Telemetry {
TracerProvider[F].tracer("org.typelevel.skunk").withVersion(BuildInfo.version).get.flatMap { implicit tracer: Tracer[F] =>
for {
operationDuration <- DbMetrics.ClientOperationDuration.create[F, Double](opDurationBoundaries)
} yield new Impl(config, connection, operationDuration)
pool <- PoolTelemetry.create[F](config.pool)
} yield new Impl(config, connection, operationDuration, pool)
}
}

def noop[F[_]: Applicative]: Telemetry[F] =
new Noop[F]

private[skunk] final class Impl[F[_]: Tracer: MonadCancelThrow](
config: TelemetryConfig,
connection: ConnectionInfo,
operationDuration: Histogram[F, Double]
operationDuration: Histogram[F, Double],
val pool: PoolTelemetry[F]
) extends Telemetry[F] {

def withConnection(connection: ConnectionInfo): Telemetry[F] =
new Impl(config, connection, operationDuration)
new Impl(config, connection, operationDuration, pool)

private val finalizationStrategy: SpanFinalizer.Strategy = {
case Resource.ExitCase.Errored(e: PostgresErrorException) =>
Expand All @@ -106,22 +111,6 @@ object Telemetry {

}

private val poolSpanF: String => F ~> F =
config.poolSpans match {
case TelemetryConfig.PoolSpans.Internal =>
label =>
FunctionK.liftFunction[F, F](
Tracer[F]
.spanBuilder(label)
.withSpanKind(SpanKind.Internal)
.build
.surround
)

case TelemetryConfig.PoolSpans.Disabled =>
Function.const(FunctionK.id[F])(_)
}

private val internalSpanF: String => F ~> F =
config.protocolSpans match {
case TelemetryConfig.ProtocolSpans.Internal =>
Expand All @@ -138,9 +127,6 @@ object Telemetry {
Function.const(FunctionK.id[F])(_)
}

def poolSpan[A](name: String)(fa: F[A]): F[A] =
poolSpanF(name)(fa)

def internalSpan[A](label: String)(fa: F[A]): F[A] =
internalSpanF(label)(fa)

Expand Down Expand Up @@ -209,6 +195,29 @@ object Telemetry {
}
}

private[skunk] final class Noop[F[_]: Applicative] extends Telemetry[F] {
private[skunk] val pool: PoolTelemetry[F] = PoolTelemetry.noop

private[skunk] def withConnection(connection: ConnectionInfo): Telemetry[F] =
this

private[skunk] def internalSpan[A](label: String)(fa: F[A]): F[A] =
fa

private[skunk] def databaseSpan[A](
operationName: String,
statement: Statement[_],
arguments: List[Option[Encoded]],
redactionStrategy: RedactionStrategy
)(fa: F[A]): F[A] = fa

private[skunk] def addAttributes(attributes: Attribute[_]*): F[Unit] =
Applicative[F].unit

private[skunk] def addProtocolAttributes(attributes: Attribute[_]*): F[Unit] =
Applicative[F].unit
}

private[skunk] def resolveOperation(
operationName: String,
statement: Statement[_],
Expand Down
35 changes: 12 additions & 23 deletions modules/core/shared/src/main/scala/telemetry/TelemetryConfig.scala
Original file line number Diff line number Diff line change
Expand Up @@ -11,8 +11,8 @@ package skunk.telemetry
sealed trait TelemetryConfig {
def captureQuery: QueryCaptureConfig
def queryAnalyzer: QueryAnalyzer
def poolSpans: TelemetryConfig.PoolSpans
def protocolSpans: TelemetryConfig.ProtocolSpans
def pool: PoolTelemetry.Config

/** Changes query text and parameter capture. */
def withCaptureQuery(captureQuery: QueryCaptureConfig): TelemetryConfig
Expand All @@ -21,25 +21,14 @@ sealed trait TelemetryConfig {
def withQueryAnalyzer(queryAnalyzer: QueryAnalyzer): TelemetryConfig

/** Enables or disables connection-pool spans. */
def withPoolSpans(poolSpans: TelemetryConfig.PoolSpans): TelemetryConfig
def withPoolConfig(config: PoolTelemetry.Config): TelemetryConfig

/** Enables or disables PostgreSQL wire-protocol spans. */
def withProtocolSpans(protocolSpans: TelemetryConfig.ProtocolSpans): TelemetryConfig
}

object TelemetryConfig {

/** Controls spans emitted by Skunk's connection pool. */
sealed trait PoolSpans
object PoolSpans {

/** Emit connection-pool operations as `INTERNAL` spans. */
case object Internal extends PoolSpans

/** Do not export connection-pool spans. */
case object Disabled extends PoolSpans
}

/** Controls spans emitted for PostgreSQL wire-protocol operations. */
sealed trait ProtocolSpans
object ProtocolSpans {
Expand All @@ -57,38 +46,38 @@ object TelemetryConfig {
val default: TelemetryConfig = TelemetryConfig(
QueryCaptureConfig.recommended,
QueryAnalyzer.noop,
PoolSpans.Disabled,
ProtocolSpans.Internal
ProtocolSpans.Internal,
PoolTelemetry.Config.default
)

/** Creates a telemetry configuration with explicit settings. */
def apply(
captureQuery: QueryCaptureConfig,
queryAnalyzer: QueryAnalyzer,
poolSpans: PoolSpans,
protocolSpans: ProtocolSpans
protocolSpans: ProtocolSpans,
pool: PoolTelemetry.Config
): TelemetryConfig =
Impl(captureQuery, queryAnalyzer, poolSpans, protocolSpans)
Impl(captureQuery, queryAnalyzer, protocolSpans, pool)

private final case class Impl(
captureQuery: QueryCaptureConfig,
queryAnalyzer: QueryAnalyzer,
poolSpans: PoolSpans,
protocolSpans: ProtocolSpans
protocolSpans: ProtocolSpans,
pool: PoolTelemetry.Config
) extends TelemetryConfig {
def withCaptureQuery(captureQuery: QueryCaptureConfig): TelemetryConfig =
copy(captureQuery = captureQuery)

def withQueryAnalyzer(queryAnalyzer: QueryAnalyzer): TelemetryConfig =
copy(queryAnalyzer = queryAnalyzer)

def withPoolSpans(poolSpans: PoolSpans): TelemetryConfig =
copy(poolSpans = poolSpans)
def withPoolConfig(config: PoolTelemetry.Config): TelemetryConfig =
copy(pool = config)

def withProtocolSpans(protocolSpans: ProtocolSpans): TelemetryConfig =
copy(protocolSpans = protocolSpans)

override def toString: String =
s"TelemetryConfig($captureQuery, $queryAnalyzer, $poolSpans, $protocolSpans)"
s"TelemetryConfig($captureQuery, $queryAnalyzer, $protocolSpans, $pool)"
}
}
Loading
Loading