diff --git a/build.sbt b/build.sbt index 293ab0e0..206db96e 100644 --- a/build.sbt +++ b/build.sbt @@ -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 += { @@ -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)) }, diff --git a/modules/core/shared/src/main/scala/Session.scala b/modules/core/shared/src/main/scala/Session.scala index fc4c865a..bee896cb 100644 --- a/modules/core/shared/src/main/scala/Session.scala +++ b/modules/core/shared/src/main/scala/Session.scala @@ -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 @@ -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 } diff --git a/modules/core/shared/src/main/scala/telemetry/ConnectionInfo.scala b/modules/core/shared/src/main/scala/telemetry/ConnectionInfo.scala index 8f65c583..ca643b57 100644 --- a/modules/core/shared/src/main/scala/telemetry/ConnectionInfo.scala +++ b/modules/core/shared/src/main/scala/telemetry/ConnectionInfo.scala @@ -4,8 +4,24 @@ package skunk.telemetry -private[skunk] final case class ConnectionInfo( - database: String, - serverAddress: String, - serverPort: Option[Long] -) +sealed trait ConnectionInfo { + 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 +} diff --git a/modules/core/shared/src/main/scala/telemetry/PoolTelemetry.scala b/modules/core/shared/src/main/scala/telemetry/PoolTelemetry.scala new file mode 100644 index 00000000..b0123f8a --- /dev/null +++ b/modules/core/shared/src/main/scala/telemetry/PoolTelemetry.scala @@ -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 + } + +} diff --git a/modules/core/shared/src/main/scala/telemetry/Telemetry.scala b/modules/core/shared/src/main/scala/telemetry/Telemetry.scala index bb436bc1..a34a2596 100644 --- a/modules/core/shared/src/main/scala/telemetry/Telemetry.scala +++ b/modules/core/shared/src/main/scala/telemetry/Telemetry.scala @@ -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} @@ -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]( @@ -42,6 +40,8 @@ sealed trait Telemetry[F[_]] { private[skunk] def addProtocolAttributes(attributes: Attribute[_]*): F[Unit] + private[skunk] def pool: PoolTelemetry[F] + } object Telemetry { @@ -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) => @@ -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 => @@ -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) @@ -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[_], diff --git a/modules/core/shared/src/main/scala/telemetry/TelemetryConfig.scala b/modules/core/shared/src/main/scala/telemetry/TelemetryConfig.scala index 5bb619d1..7f044f4d 100644 --- a/modules/core/shared/src/main/scala/telemetry/TelemetryConfig.scala +++ b/modules/core/shared/src/main/scala/telemetry/TelemetryConfig.scala @@ -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 @@ -21,7 +21,7 @@ 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 @@ -29,17 +29,6 @@ sealed trait 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 { @@ -57,24 +46,24 @@ 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) @@ -82,13 +71,13 @@ object TelemetryConfig { 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)" } } diff --git a/modules/core/shared/src/main/scala/util/Pool.scala b/modules/core/shared/src/main/scala/util/Pool.scala index 968ea63a..7812c01b 100644 --- a/modules/core/shared/src/main/scala/util/Pool.scala +++ b/modules/core/shared/src/main/scala/util/Pool.scala @@ -12,7 +12,7 @@ import cats.effect.implicits._ import cats.effect.Resource import cats.syntax.all._ import skunk.exception.SkunkException -import skunk.telemetry.Telemetry +import skunk.telemetry.PoolTelemetry object Pool { @@ -43,27 +43,19 @@ object Pool { """.stripMargin.trim.linesIterator.mkString(" ")) ) - // Preserved for previous use, and specifically simpler use for - // Tracer systems that are universal rather than shorter scoped. - def of[F[_]: Concurrent: Telemetry, A]( - rsrc: Resource[F, A], - size: Int)( - recycler: Recycler[F, A] - ): Resource[F, Resource[F, A]] = ofF({(_: Telemetry[F]) => rsrc}, size)(recycler).map(_.apply(Telemetry[F])) - /** - * A pooled resource (which is itself a managed resource). + * Preserved for previous use, and specifically simpler use for + * Tracer systems that are universal rather than shorter scoped. * @param rsrc the underlying resource to be pooled * @param size maximum size of the pool (must be positive) * @param recycler a cleanup/health-check to be done before elements are returned to the pool; * yielding false here means the element should be freed and removed from the pool. */ - def ofF[F[_]: Concurrent, A]( - rsrc: Telemetry[F] => Resource[F, A], + def of[F[_]: Concurrent: PoolTelemetry, A]( + rsrc: Resource[F, A], size: Int)( recycler: Recycler[F, A] - ): Resource[F, Telemetry[F] => Resource[F, A]] = - ofF(rsrc, size, checkout = Recycler.success[F, A], checkin = recycler) + ): Resource[F, Resource[F, A]] = of(rsrc, size, checkout = Recycler.success[F, A], checkin = recycler) /** * A pooled resource (which is itself a managed resource). @@ -74,12 +66,12 @@ object Pool { * @param checkin a cleanup/health-check to be done before elements are returned to the pool; * yielding false here means the element should be freed and removed from the pool. */ - def ofF[F[_]: Concurrent, A]( - rsrc: Telemetry[F] => Resource[F, A], + def of[F[_]: Concurrent: PoolTelemetry, A]( + rsrc: Resource[F, A], size: Int, checkout: Recycler[F, A], checkin: Recycler[F, A] - ): Resource[F, Telemetry[F] => Resource[F, A]] = { + ): Resource[F, Resource[F, A]] = { // Just in case. assert(size > 0, s"Pool size must be positive (you passed $size).") @@ -95,14 +87,14 @@ object Pool { ) // We can construct a pool given a Ref containing our initial state. - def poolImpl(ref: Ref[F, State])(implicit T: Telemetry[F]): Resource[F, A] = { + def poolImpl(ref: Ref[F, State]): Resource[F, A] = { // To give out an alloc we create a deferral first, which we will need if there are no slots // available. If there is a filled slot, remove it and yield its alloc. If there is an empty // slot, remove it and allocate. If there are no slots, enqueue the deferral and wait on it, // which will [semantically] block the caller until an alloc is returned to the pool. def give(poll: Poll[F]): F[Alloc] = - Telemetry[F].poolSpan("pool.allocate")(giveLoop(poll)) + PoolTelemetry[F].span("pool.allocate")(giveLoop(poll)) def giveLoop(poll: Poll[F]): F[Alloc] = Deferred[F, Either[Throwable, Alloc]].flatMap { d => @@ -114,36 +106,11 @@ object Pool { case _ => ref.update { case (os, ds) => (os :+ None, ds) } } - // Here we go. The cases are a full slot (done), an empty slot (alloc), and no slots at - // all (defer and wait). - ref.modify { - case (Some(a) :: os, ds) => ((os, ds), a.pure[F]) - case (None :: os, ds) => ((os, ds), Concurrent[F].onError(rsrc(Telemetry[F]).allocated)(restore)) - case (Nil, ds) => - val cancel = ref.flatModify { // try to remove our deferred - case (os, ds) => - val canRemove = ds.contains(d) - val cleanupMaybe = if (canRemove) // we'll pull it out before anyone can complete it - ().pure[F] - else // someone got to it first and will complete it, so we wait and then return it - d.get.flatMap(_.liftTo[F]).onError(restore).flatMap(take(_)) - - ((os, if (canRemove) ds.filterNot(_ == d) else ds), cleanupMaybe) - } - - val wait = - poll(d.get) - .onCancel(cancel) - .flatMap(_.liftTo[F].onError(restore)) - ((Nil, ds :+ d), wait) - } .flatten - - // Here we go. The cases are a full slot (check and done), an empty slot (alloc), and no // slots at all (defer and wait). ref.modify { case (Some(a) :: os, ds) => ((os, ds), reuse(poll, a)) - case (None :: os, ds) => ((os, ds), Concurrent[F].onError(rsrc(Telemetry[F]).allocated)(restore)) + case (None :: os, ds) => ((os, ds), Concurrent[F].onError(rsrc.allocated)(restore)) case (Nil, ds) => val cancel = ref.flatModify { // try to remove our deferred case (os, ds) => @@ -181,7 +148,7 @@ object Pool { // there are a bunch of error conditions to consider. This operation is a finalizer and // cannot be canceled, so we don't need to worry about that case here. def take(a: Alloc): F[Unit] = - Telemetry[F].poolSpan("pool.free") { + PoolTelemetry[F].span("pool.free") { checkin(a._1).onError { case _ => dispose(a) } flatMap { case true => recycle(a) case false => dispose(a) @@ -191,7 +158,7 @@ object Pool { // Return `a` to the pool. If there are awaiting deferrals, complete the next one. Otherwise // push a filled slot into the queue. def recycle(a: Alloc): F[Unit] = - Telemetry[F].poolSpan("recycle") { + PoolTelemetry[F].span("recycle") { ref.modify { case (os, d :: ds) => ((os, ds), d.complete(a.asRight).void) // hand it back out case (os, Nil) => ((Some(a) :: os, Nil), ().pure[F]) // return to pool @@ -203,10 +170,10 @@ object Pool { // of `a`. If there are deferrals, remove the next one and complete it (failures in allocation // are handled by the awaiting deferral in `give` above). Always finalize `a` def dispose(a: Alloc): F[Unit] = - Telemetry[F].poolSpan("dispose") { + PoolTelemetry[F].span("dispose") { ref.modify { case (os, Nil) => ((os :+ None, Nil), ().pure[F]) // new empty slot - case (os, d :: ds) => ((os, ds), Concurrent[F].attempt(rsrc(Telemetry[F]).allocated).flatMap(d.complete).void) // alloc now! + case (os, d :: ds) => ((os, ds), Concurrent[F].attempt(rsrc.allocated).flatMap(d.complete).void) // alloc now! }.flatMap(next => a._2.guarantee(next)) // first finalize the original alloc then potentially do new alloc } @@ -239,7 +206,7 @@ object Pool { } - Resource.make(alloc)(free).map(a => {implicit T: Telemetry[F] => poolImpl(a)}) + Resource.make(alloc)(free).map(a => poolImpl(a)) } diff --git a/modules/docs/src/main/laika/tutorial/Telemetry.md b/modules/docs/src/main/laika/tutorial/Telemetry.md index ac958d91..f81adc8d 100644 --- a/modules/docs/src/main/laika/tutorial/Telemetry.md +++ b/modules/docs/src/main/laika/tutorial/Telemetry.md @@ -108,12 +108,16 @@ val protocolTelemetry = TelemetryConfig.default .withProtocolSpans(TelemetryConfig.ProtocolSpans.Disabled) ``` -Pool spans are disabled by default. They can be enabled when investigating connection acquisition -or pool cleanup: +Pool spans are disabled by default through `PoolTelemetry.Config.default`. Configure them with +`TelemetryConfig.withPoolConfig` when investigating connection acquisition or pool cleanup: ```scala mdoc:silent +import skunk.telemetry.PoolTelemetry + val poolTelemetry = TelemetryConfig.default - .withPoolSpans(TelemetryConfig.PoolSpans.Internal) + .withPoolConfig( + PoolTelemetry.Config.default.withPoolSpans(PoolTelemetry.Config.PoolSpans.Internal) + ) ``` To state explicitly that neither category should be emitted: @@ -121,7 +125,9 @@ To state explicitly that neither category should be emitted: ```scala mdoc:silent val minimalTelemetry = TelemetryConfig.default .withProtocolSpans(TelemetryConfig.ProtocolSpans.Disabled) - .withPoolSpans(TelemetryConfig.PoolSpans.Disabled) + .withPoolConfig( + PoolTelemetry.Config.default.withPoolSpans(PoolTelemetry.Config.PoolSpans.Disabled) + ) val sessions = Session.Builder[IO] diff --git a/modules/tests/shared/src/main/scala/ffstest/FFramework.scala b/modules/tests/shared/src/main/scala/ffstest/FFramework.scala index 3f0f8877..5a6df348 100644 --- a/modules/tests/shared/src/main/scala/ffstest/FFramework.scala +++ b/modules/tests/shared/src/main/scala/ffstest/FFramework.scala @@ -40,12 +40,6 @@ trait FTest extends CatsEffectSuite with FTestPlatform { def tracedTest[A](options: TestOptions)(body: TracerProvider[IO] => IO[A])(implicit loc: Location): Unit = test(options)(withinSpan(options.name)((provider, _) => body(provider))) - def tracedTestWithTracer[A](name: String)(body: Tracer[IO] => IO[A])(implicit loc: Location): Unit = - test(name)(withinSpan(name)((_, tracer) => body(tracer))) - - def tracedTestWithTracer[A](options: TestOptions)(body: Tracer[IO] => IO[A])(implicit loc: Location): Unit = - test(options)(withinSpan(options.name)((_, tracer) => body(tracer))) - def pureTest(name: String)(f: => Boolean): Unit = test(name)(assert(name, f)) def fail[A](msg: String): IO[A] = IO.raiseError(new AssertionError(msg)) def fail[A](msg: String, cause: Throwable): IO[A] = IO.raiseError(new AssertionError(msg, cause)) diff --git a/modules/tests/jvm/src/test/scala/TelemetryPoolConfigTest.scala b/modules/tests/shared/src/test/scala/PoolTelemetryConfigTest.scala similarity index 73% rename from modules/tests/jvm/src/test/scala/TelemetryPoolConfigTest.scala rename to modules/tests/shared/src/test/scala/PoolTelemetryConfigTest.scala index f4ad140f..b7e8940e 100644 --- a/modules/tests/jvm/src/test/scala/TelemetryPoolConfigTest.scala +++ b/modules/tests/shared/src/test/scala/PoolTelemetryConfigTest.scala @@ -6,7 +6,6 @@ package skunk import cats.effect.IO import munit.CatsEffectSuite -import org.typelevel.otel4s.metrics.MeterProvider import org.typelevel.otel4s.sdk.testkit.InstrumentationScopeExpectation import org.typelevel.otel4s.sdk.testkit.OpenTelemetrySdkTestkit import org.typelevel.otel4s.sdk.testkit.trace.SpanExpectation @@ -16,11 +15,9 @@ import org.typelevel.otel4s.sdk.testkit.trace.TraceExpectations import org.typelevel.otel4s.sdk.testkit.trace.TraceForestExpectation import org.typelevel.otel4s.sdk.trace.data.SpanData import org.typelevel.otel4s.trace.TracerProvider -import skunk.telemetry.ConnectionInfo -import skunk.telemetry.Telemetry -import skunk.telemetry.TelemetryConfig +import skunk.telemetry.PoolTelemetry -class TelemetryPoolConfigTest extends CatsEffectSuite { +class PoolTelemetryConfigTest extends CatsEffectSuite { private val scope = InstrumentationScopeExpectation @@ -28,24 +25,23 @@ class TelemetryPoolConfigTest extends CatsEffectSuite { .version(BuildInfo.version) .attributesEmpty - private def poolSpans(config: TelemetryConfig): IO[List[SpanData]] = + private def poolSpans(config: PoolTelemetry.Config): IO[List[SpanData]] = OpenTelemetrySdkTestkit.inMemory[IO]().use { testkit => implicit val tracerProvider: TracerProvider[IO] = testkit.tracerProvider - implicit val meterProvider: MeterProvider[IO] = testkit.meterProvider - Telemetry - .create[IO](config, ConnectionInfo("world", "localhost", None)) - .flatMap(_.poolSpan("pool.allocate")(IO.unit)) *> + PoolTelemetry + .create[IO](config) + .flatMap(_.span("pool.allocate")(IO.unit)) *> testkit.finishedSpans } test("pool spans are disabled by default") { - poolSpans(TelemetryConfig.default).map(spans => assertEquals(spans, Nil)) + poolSpans(PoolTelemetry.Config.default).map(spans => assertEquals(spans, Nil)) } test("pool spans can be emitted as internal spans") { poolSpans( - TelemetryConfig.default.withPoolSpans(TelemetryConfig.PoolSpans.Internal) + PoolTelemetry.Config(PoolTelemetry.Config.PoolSpans.Internal) ).map { spans => val expectation = TraceForestExpectation.unordered( diff --git a/modules/tests/shared/src/test/scala/PoolTest.scala b/modules/tests/shared/src/test/scala/PoolTest.scala index d25ab9e9..4dd9157f 100644 --- a/modules/tests/shared/src/test/scala/PoolTest.scala +++ b/modules/tests/shared/src/test/scala/PoolTest.scala @@ -15,16 +15,13 @@ import skunk.util.Pool.ResourceLeak import cats.effect.Deferred import scala.util.Random import skunk.util.Pool.ShutdownException -import org.typelevel.otel4s.trace.Tracer import skunk.util.Recycler +import skunk.telemetry.PoolTelemetry import cats.effect.testkit.TestControl -import skunk.telemetry.Telemetry +import munit.Location class PoolTest extends FTest { - implicit def telemetry(implicit tracer: Tracer[IO]): Telemetry[IO] = - skunk.TestTelemetry("pool-test") - case class UserFailure() extends Exception("user failure") case class AllocFailure() extends Exception("allocation failure") case class FreeFailure() extends Exception("free failure") @@ -36,6 +33,9 @@ class PoolTest extends FTest { Resource.make(next)(_ => IO.unit) } + val poolTelemetryConfig = + PoolTelemetry.Config(PoolTelemetry.Config.PoolSpans.Internal) + // list of computations into computation that yields results one by one def yielding[A](fas: IO[A]*): IO[IO[A]] = Ref[IO].of(fas.toList).map { ref => @@ -48,21 +48,26 @@ class PoolTest extends FTest { def resourceYielding[A](fas: IO[A]*): IO[Resource[IO, A]] = yielding(fas: _*).map(Resource.make(_)(_ => IO.unit)) + def poolTelemetryTest[A](name: String)(body: PoolTelemetry[IO] => IO[A])(implicit loc: Location): Unit = + tracedTest(name) { implicit tracerProvider => + PoolTelemetry.create[IO](poolTelemetryConfig).flatMap(telemetry => body(telemetry)) + } + // This test leaks - tracedTestWithTracer("error in alloc is rethrown to caller (immediate)") { implicit tracer: Tracer[IO] => + poolTelemetryTest("error in alloc is rethrown to caller (immediate)") { implicit telemetry => val rsrc = Resource.make(IO.raiseError[String](AllocFailure()))(_ => IO.unit) - val pool = Pool.ofF({(_: Telemetry[IO]) => rsrc}, 42)(Recycler.success) - pool.use(_(Telemetry[IO]).use(_ => IO.unit)).assertFailsWith[AllocFailure] + val pool = Pool.of(rsrc, 42)(Recycler.success) + pool.use(_.use(_ => IO.unit)).assertFailsWith[AllocFailure] } - tracedTestWithTracer("error in alloc is rethrown to caller (deferral completion following errored cleanup)") { implicit tracer: Tracer[IO] => + poolTelemetryTest("error in alloc is rethrown to caller (deferral completion following errored cleanup)") { implicit telemetry => resourceYielding(IO(1), IO.raiseError(AllocFailure())).flatMap { r => - val p = Pool.ofF({(_: Telemetry[IO]) => r}, 1)(Recycler[IO, Int](_ => IO.raiseError(ResetFailure()))) + val p = Pool.of(r, 1)(Recycler[IO, Int](_ => IO.raiseError(ResetFailure()))) p.use { r => for { d <- Deferred[IO, Unit] - f1 <- r(Telemetry[IO]).use(n => assertEqual("n should be 1", n, 1) *> d.get).assertFailsWith[ResetFailure].start - f2 <- r(Telemetry[IO]).use(_ => fail[Int]("should never get here")).assertFailsWith[AllocFailure].start + f1 <- r.use(n => assertEqual("n should be 1", n, 1) *> d.get).assertFailsWith[ResetFailure].start + f2 <- r.use(_ => fail[Int]("should never get here")).assertFailsWith[AllocFailure].start _ <- d.complete(()) _ <- f1.join _ <- f2.join @@ -71,14 +76,14 @@ class PoolTest extends FTest { } } - tracedTestWithTracer("error in alloc is rethrown to caller (deferral completion following failed cleanup)") { implicit tracer: Tracer[IO] => + poolTelemetryTest("error in alloc is rethrown to caller (deferral completion following failed cleanup)") { implicit telemetry => resourceYielding(IO(1), IO.raiseError(AllocFailure())).flatMap { r => - val p = Pool.ofF({(_: Telemetry[IO]) => r}, 1)(Recycler.failure) + val p = Pool.of(r, 1)(Recycler.failure) p.use { r => for { d <- Deferred[IO, Unit] - f1 <- r(Telemetry[IO]).use(n => assertEqual("n should be 1", n, 1) *> d.get).start - f2 <- r(Telemetry[IO]).use(_ => fail[Int]("should never get here")).assertFailsWith[AllocFailure].start + f1 <- r.use(n => assertEqual("n should be 1", n, 1) *> d.get).start + f2 <- r.use(_ => fail[Int]("should never get here")).assertFailsWith[AllocFailure].start _ <- d.complete(()) _ <- f1.join _ <- f2.join @@ -87,25 +92,25 @@ class PoolTest extends FTest { } } - tracedTestWithTracer("error in finalizer does not prevent cleanup of deferreds") { implicit tracer: Tracer[IO] => + poolTelemetryTest("error in finalizer does not prevent cleanup of deferreds") { implicit telemetry => val r = Resource.make(IO(1))(_ => IO.raiseError(ResetFailure())) - val p = Pool.ofF({(_: Telemetry[IO]) => r}, 1)(Recycler.failure) + val p = Pool.of(r, 1)(Recycler.failure) p.use { r => - val tx = r(Telemetry[IO]).use(_ => IO.unit) + val tx = r.use(_ => IO.unit) List(tx, tx).parSequence }.assertFailsWith[ResetFailure] } - tracedTestWithTracer("provoke dangling deferral cancellation") { implicit tracer: Tracer[IO] => + poolTelemetryTest("provoke dangling deferral cancellation") { implicit telemetry => ints.flatMap { r => - val p = Pool.ofF({(_: Telemetry[IO]) => r}, 1)(Recycler.failure) + val p = Pool.of(r, 1)(Recycler.failure) Deferred[IO, Either[Throwable, Int]].flatMap { d1 => p.use { r => for { d <- Deferred[IO, Unit] - _ <- r(Telemetry[IO]).use(_ => d.complete(()) *> IO.never).start // leaked forever + _ <- r.use(_ => d.complete(()) *> IO.never).start // leaked forever _ <- d.get // make sure the resource has been allocated - f <- r(Telemetry[IO]).use(_ => fail[Int]("should never get here")).attempt.flatMap(d1.complete).start // defer + f <- r.use(_ => fail[Int]("should never get here")).attempt.flatMap(d1.complete).start // defer _ <- IO.sleep(100.milli) // ensure that the fiber has a chance to run } yield f } .assertFailsWith[ResourceLeak].flatMap { @@ -115,40 +120,40 @@ class PoolTest extends FTest { } }} - tracedTestWithTracer("error in free is rethrown to caller") { implicit tracer: Tracer[IO] => + poolTelemetryTest("error in free is rethrown to caller") { implicit telemetry => val rsrc = Resource.make("foo".pure[IO])(_ => IO.raiseError(FreeFailure())) - val pool = Pool.ofF({(_: Telemetry[IO]) => rsrc}, 42)(Recycler.success) - pool.use(_(Telemetry[IO]).use(_ => IO.unit)).assertFailsWith[FreeFailure] + val pool = Pool.of(rsrc, 42)(Recycler.success) + pool.use(_.use(_ => IO.unit)).assertFailsWith[FreeFailure] } - tracedTestWithTracer("error in reset is rethrown to caller") { implicit tracer: Tracer[IO] => + poolTelemetryTest("error in reset is rethrown to caller") { implicit telemetry => val rsrc = Resource.make("foo".pure[IO])(_ => IO.unit) - val pool = Pool.ofF({(_: Telemetry[IO]) => rsrc}, 42)(Recycler[IO, String](_ => IO.raiseError(ResetFailure()))) - pool.use(_(Telemetry[IO]).use(_ => IO.unit)).assertFailsWith[ResetFailure] + val pool = Pool.of(rsrc, 42)(Recycler[IO, String](_ => IO.raiseError(ResetFailure()))) + pool.use(_.use(_ => IO.unit)).assertFailsWith[ResetFailure] } - tracedTestWithTracer("reuse on serial access") { implicit tracer: Tracer[IO] => - ints.map(a => Pool.ofF({(_: Telemetry[IO]) => a}, 3)(Recycler.success)).flatMap { factory => + poolTelemetryTest("reuse on serial access") { implicit telemetry => + ints.map(a => Pool.of(a, 3)(Recycler.success)).flatMap { factory => factory.use { pool => - pool(Telemetry[IO]).use { n => + pool.use { n => assertEqual("first num should be 1", n, 1) } *> - pool(Telemetry[IO]).use { n => + pool.use { n => assertEqual("we should get it again", n, 1) } } } } - tracedTestWithTracer("allocation on nested access") { implicit tracer: Tracer[IO] => - ints.map(a => Pool.ofF({(_: Telemetry[IO]) => a}, 3)(Recycler.success)).flatMap { factory => + poolTelemetryTest("allocation on nested access") { implicit telemetry => + ints.map(a => Pool.of(a, 3)(Recycler.success)).flatMap { factory => factory.use { pool => - pool(Telemetry[IO]).use { n => + pool.use { n => assertEqual("first num should be 1", n, 1) *> - pool(Telemetry[IO]).use { n => + pool.use { n => assertEqual("but this one should be 2", n, 2) } *> - pool(Telemetry[IO]).use { n => + pool.use { n => assertEqual("and again", n, 2) } } @@ -156,10 +161,10 @@ class PoolTest extends FTest { } } - tracedTestWithTracer("allocated resource can cause a leak, which will be detected on finalization") { implicit tracer: Tracer[IO] => - ints.map(a => Pool.ofF({(_: Telemetry[IO]) => a}, 3)(Recycler.success)).flatMap { factory => + poolTelemetryTest("allocated resource can cause a leak, which will be detected on finalization") { implicit telemetry => + ints.map(a => Pool.of(a, 3)(Recycler.success)).flatMap { factory => factory.use { pool => - pool(Telemetry[IO]).allocated + pool.allocated } .assertFailsWith[ResourceLeak].flatMap { case ResourceLeak(expected, actual, _) => assert("expected 1 leakage", expected - actual == 1) @@ -167,10 +172,10 @@ class PoolTest extends FTest { } } - tracedTestWithTracer("unmoored fiber can cause a leak, which will be detected on finalization") { implicit tracer: Tracer[IO] => - ints.map(a => Pool.ofF({(_: Telemetry[IO]) => a}, 3)(Recycler.success)).flatMap { factory => + poolTelemetryTest("unmoored fiber can cause a leak, which will be detected on finalization") { implicit telemetry => + ints.map(a => Pool.of(a, 3)(Recycler.success)).flatMap { factory => factory.use { pool => - pool(Telemetry[IO]).use(_ => IO.never).start *> + pool.use(_ => IO.never).start *> IO.sleep(100.milli) // ensure that the fiber has a chance to run } .assertFailsWith[ResourceLeak].flatMap { case ResourceLeak(expected, actual, _) => @@ -186,11 +191,11 @@ class PoolTest extends FTest { val shortRandomDelay = IO((Random.nextInt() % 100).abs.milliseconds) - tracedTestWithTracer("progress and safety with many fibers") { implicit tracer: Tracer[IO] => - ints.map(a => Pool.ofF({(_: Telemetry[IO]) => a}, PoolSize)(Recycler.success)).flatMap { factory => + poolTelemetryTest("progress and safety with many fibers") { implicit telemetry => + ints.map(a => Pool.of(a, PoolSize)(Recycler.success)).flatMap { factory => (1 to ConcurrentTasks).toList.parTraverse_{ _ => factory.use { p => - p(Telemetry[IO]).use { _ => + p.use { _ => for { t <- shortRandomDelay _ <- IO.sleep(t) @@ -201,13 +206,13 @@ class PoolTest extends FTest { } } - tracedTestWithTracer("progress and safety with many fibers and cancellation") { implicit tracer: Tracer[IO] => - ints.map(a => Pool.ofF({(_: Telemetry[IO]) => a}, PoolSize)(Recycler.success)).flatMap { factory => + poolTelemetryTest("progress and safety with many fibers and cancellation") { implicit telemetry => + ints.map(a => Pool.of(a, PoolSize)(Recycler.success)).flatMap { factory => factory.use { pool => (1 to ConcurrentTasks).toList.parTraverse_{_ => for { t <- shortRandomDelay - f <- pool(Telemetry[IO]).use(_ => IO.sleep(t)).start + f <- pool.use(_ => IO.sleep(t)).start _ <- if (t > 50.milliseconds) f.join else f.cancel } yield () } @@ -215,11 +220,11 @@ class PoolTest extends FTest { } } - tracedTestWithTracer("progress and safety with many fibers and user failures") { implicit tracer: Tracer[IO] => - ints.map(a => Pool.ofF({(_: Telemetry[IO]) => a}, PoolSize)(Recycler.success)).flatMap { factory => + poolTelemetryTest("progress and safety with many fibers and user failures") { implicit telemetry => + ints.map(a => Pool.of(a, PoolSize)(Recycler.success)).flatMap { factory => factory.use { pool => (1 to ConcurrentTasks).toList.parTraverse_{ _ => - pool(Telemetry[IO]).use { _ => + pool.use { _ => for { t <- shortRandomDelay _ <- IO.sleep(t) @@ -231,30 +236,30 @@ class PoolTest extends FTest { } } - tracedTestWithTracer("progress and safety with many fibers and allocation failures") { implicit tracer: Tracer[IO] => + poolTelemetryTest("progress and safety with many fibers and allocation failures") { implicit telemetry => val alloc = IO(Random.nextBoolean()).flatMap { case true => IO.unit case false => IO.raiseError(AllocFailure()) } val rsrc = Resource.make(alloc)(_ => IO.unit) - Pool.ofF({(_: Telemetry[IO]) => rsrc}, PoolSize)(Recycler.success).use { pool => + Pool.of(rsrc, PoolSize)(Recycler.success).use { pool => (1 to ConcurrentTasks).toList.parTraverse_{ _ => - pool(Telemetry[IO]).use { _ => + pool.use { _ => IO.unit } .attempt } } } - tracedTestWithTracer("progress and safety with many fibers and freeing failures") { implicit tracer: Tracer[IO] => + poolTelemetryTest("progress and safety with many fibers and freeing failures") { implicit telemetry => val free = IO(Random.nextBoolean()).flatMap { case true => IO.unit case false => IO.raiseError(FreeFailure()) } val rsrc = Resource.make(IO.unit)(_ => free) - Pool.ofF({(_: Telemetry[IO]) => rsrc}, PoolSize)(Recycler.success).use { pool => + Pool.of(rsrc, PoolSize)(Recycler.success).use { pool => (1 to ConcurrentTasks).toList.parTraverse_{ _ => - pool(Telemetry[IO]).use { _ => + pool.use { _ => IO.unit } .attempt } @@ -265,16 +270,16 @@ class PoolTest extends FTest { } } - tracedTestWithTracer("progress and safety with many fibers and reset failures") { implicit tracer: Tracer[IO] => + poolTelemetryTest("progress and safety with many fibers and reset failures") { implicit telemetry => val recycle = IO(Random.nextInt(3)).flatMap { case 0 => true.pure[IO] case 1 => false.pure[IO] case 2 => IO.raiseError(ResetFailure()) } val rsrc = Resource.make(IO.unit)(_ => IO.unit) - Pool.ofF({(_: Telemetry[IO]) => rsrc}, PoolSize)(Recycler(_ => recycle)).use { pool => + Pool.of(rsrc, PoolSize)(Recycler(_ => recycle)).use { pool => (1 to ConcurrentTasks).toList.parTraverse_{ _ => - pool(Telemetry[IO]).use { _ => + pool.use { _ => IO.unit } handleErrorWith { case ResetFailure() => IO.unit @@ -284,37 +289,37 @@ class PoolTest extends FTest { } } - tracedTestWithTracer("unhealthy pooled resource is discarded and replaced on checkout") { implicit tracer: Tracer[IO] => + poolTelemetryTest("unhealthy pooled resource is discarded and replaced on checkout") { implicit telemetry => for { healthy <- Ref[IO].of(true) freed <- Ref[IO].of(List.empty[Int]) counter <- Ref[IO].of(1) rsrc = Resource.make(counter.getAndUpdate(_ + 1))(n => freed.update(_ :+ n)) - _ <- Pool.ofF({(_: Telemetry[IO]) => rsrc}, 1, Recycler[IO, Int](_ => healthy.get), Recycler.success[IO, Int]).use { pool => + _ <- Pool.of(rsrc, 1, Recycler[IO, Int](_ => healthy.get), Recycler.success[IO, Int]).use { pool => for { - _ <- pool(Telemetry[IO]).use(n => assertEqual("first checkout", n, 1)) + _ <- pool.use(n => assertEqual("first checkout", n, 1)) _ <- healthy.set(false) - _ <- pool(Telemetry[IO]).use(n => assertEqual("second checkout", n, 2)) + _ <- pool.use(n => assertEqual("second checkout", n, 2)) _ <- freed.get.flatMap(assertEqual("the dead resource was freed", _, List(1))) } yield () } } yield () } - tracedTestWithTracer("health check that raises discards the resource rather than failing the checkout") { implicit tracer: Tracer[IO] => + poolTelemetryTest("health check that raises discards the resource rather than failing the checkout") { implicit telemetry => for { first <- Ref[IO].of(true) counter <- Ref[IO].of(1) rsrc = Resource.make(counter.getAndUpdate(_ + 1))(_ => IO.unit) check = Recycler[IO, Int](_ => first.getAndSet(false).ifM(IO.raiseError[Boolean](AllocFailure()), true.pure[IO])) - _ <- Pool.ofF({(_: Telemetry[IO]) => rsrc}, 1, check, Recycler.success[IO, Int]).use { pool => - pool(Telemetry[IO]).use_ *> - pool(Telemetry[IO]).use(n => assertEqual("second checkout", n, 2)) + _ <- Pool.of(rsrc, 1, check, Recycler.success[IO, Int]).use { pool => + pool.use_ *> + pool.use(n => assertEqual("second checkout", n, 2)) } } yield () } - test("cancel while waiting") { implicit tracer: Tracer[IO] => + poolTelemetryTest("cancel while waiting") { implicit telemetry => TestControl.executeEmbed { Pool.of(Resource.unit[IO], 1)(Recycler.success).use { pool => pool.useForever.background.surround { // take away the resource ... diff --git a/modules/tests/shared/src/test/scala/TelemetryConfigTest.scala b/modules/tests/shared/src/test/scala/TelemetryConfigTest.scala index e88193f7..742187dd 100644 --- a/modules/tests/shared/src/test/scala/TelemetryConfigTest.scala +++ b/modules/tests/shared/src/test/scala/TelemetryConfigTest.scala @@ -12,19 +12,11 @@ import skunk.data.Encoded import skunk.telemetry.ConnectionInfo import skunk.telemetry.QueryAnalyzer import skunk.telemetry.QueryCaptureConfig +import skunk.telemetry.PoolTelemetry import skunk.telemetry.Telemetry import skunk.telemetry.TelemetryConfig import skunk.util.Origin -object TestTelemetry { - def apply(database: String)(implicit tracer: org.typelevel.otel4s.trace.Tracer[cats.effect.IO]): Telemetry[cats.effect.IO] = - new Telemetry.Impl( - TelemetryConfig.default, - ConnectionInfo(database, "simulated", None), - org.typelevel.otel4s.metrics.Histogram.noop[cats.effect.IO, Double], - ) -} - class TelemetryConfigTest extends FunSuite { private val connection = @@ -41,7 +33,7 @@ class TelemetryConfigTest extends FunSuite { val config = TelemetryConfig.default .withCaptureQuery(capture) .withQueryAnalyzer(analyzer) - .withPoolSpans(TelemetryConfig.PoolSpans.Internal) + .withPoolConfig(PoolTelemetry.Config(PoolTelemetry.Config.PoolSpans.Internal)) .withProtocolSpans(TelemetryConfig.ProtocolSpans.Disabled) val analysis = QueryAnalyzer.Analysis( queryText = Some("SELECT ?"), @@ -54,7 +46,7 @@ class TelemetryConfigTest extends FunSuite { assertEquals(capture.queryParametersPolicy, QueryCaptureConfig.QueryParametersPolicy.All) assertEquals(config.captureQuery, capture) assertEquals(config.queryAnalyzer, analyzer) - assertEquals(config.poolSpans, TelemetryConfig.PoolSpans.Internal) + assertEquals(config.pool.poolSpans, PoolTelemetry.Config.PoolSpans.Internal) assertEquals(config.protocolSpans, TelemetryConfig.ProtocolSpans.Disabled) assertEquals(analysis.queryText, Some("SELECT ?")) assertEquals(analysis.storedProcedureName, Some("find_country")) @@ -63,7 +55,7 @@ class TelemetryConfigTest extends FunSuite { } test("pool spans are disabled by default") { - assertEquals(TelemetryConfig.default.poolSpans, TelemetryConfig.PoolSpans.Disabled) + assertEquals(TelemetryConfig.default.pool.poolSpans, PoolTelemetry.Config.PoolSpans.Disabled) } test("QueryAnalyzer.noop and fallback") { @@ -145,7 +137,7 @@ class TelemetryConfigTest extends FunSuite { Nil, RedactionStrategy.OptIn, TelemetryConfig.default.withQueryAnalyzer(analyzer), - connection.copy(serverPort = Some(6432L)), + ConnectionInfo(connection.database, connection.serverAddress, Some(6432L)), ) assertEquals(operation.spanName, "get country by code") diff --git a/modules/tests/jvm/src/test/scala/TelemetryIntegrationTest.scala b/modules/tests/shared/src/test/scala/TelemetryIntegrationTest.scala similarity index 100% rename from modules/tests/jvm/src/test/scala/TelemetryIntegrationTest.scala rename to modules/tests/shared/src/test/scala/TelemetryIntegrationTest.scala diff --git a/modules/tests/jvm/src/test/scala/TelemetryQueryAnalyzerTest.scala b/modules/tests/shared/src/test/scala/TelemetryQueryAnalyzerTest.scala similarity index 100% rename from modules/tests/jvm/src/test/scala/TelemetryQueryAnalyzerTest.scala rename to modules/tests/shared/src/test/scala/TelemetryQueryAnalyzerTest.scala diff --git a/modules/tests/shared/src/test/scala/simulation/ExchangeCountTest.scala b/modules/tests/shared/src/test/scala/simulation/ExchangeCountTest.scala index 16d65ecc..db2cbc74 100644 --- a/modules/tests/shared/src/test/scala/simulation/ExchangeCountTest.scala +++ b/modules/tests/shared/src/test/scala/simulation/ExchangeCountTest.scala @@ -8,8 +8,6 @@ package simulation import cats.effect.IO import cats.syntax.all._ import ffstest.FTest -import org.typelevel.otel4s.metrics.Histogram -import org.typelevel.otel4s.trace.Tracer import skunk.{ RedactionStrategy, Session, TypingStrategy } import skunk.codec.all._ import skunk.data.{ Completion, TransactionStatus, Type } @@ -29,8 +27,7 @@ import skunk.util.{ Namer, Typer } */ class ExchangeCountTest extends FTest with SimMessageSocket.DSL { - implicit val tracer: Tracer[IO] = Tracer.noop - implicit val telemetry: Telemetry[IO] = skunk.TestTelemetry("simulated") + implicit val telemetry: Telemetry[IO] = Telemetry.noop private val int4Column: RowDescription.Field = RowDescription.Field("?column?", 0, 0, Typer.Static.oidForType(Type.int4).get, 4, 0, 0) diff --git a/modules/tests/shared/src/test/scala/simulation/SimTest.scala b/modules/tests/shared/src/test/scala/simulation/SimTest.scala index 22072d23..20be2060 100644 --- a/modules/tests/shared/src/test/scala/simulation/SimTest.scala +++ b/modules/tests/shared/src/test/scala/simulation/SimTest.scala @@ -8,7 +8,6 @@ package simulation import cats.effect._ import ffstest.FTest import fs2.concurrent.Signal -import org.typelevel.otel4s.trace.Tracer import skunk.{Session, RedactionStrategy, TypingStrategy} import skunk.data.Notification import skunk.data.TransactionStatus @@ -24,8 +23,7 @@ import skunk.telemetry.Telemetry trait SimTest extends FTest with SimMessageSocket.DSL { - implicit val tracer: Tracer[IO] = Tracer.noop - implicit val telemetry: Telemetry[IO] = skunk.TestTelemetry("simulated") + implicit val telemetry: Telemetry[IO] = Telemetry.noop private class SimulatedBufferedMessageSocket(ms: MessageSocket[IO]) extends BufferedMessageSocket[IO] { def receive: IO[BackendMessage] = ms.receive