diff --git a/.devcontainer/db/init-db.sh b/.devcontainer/db/init-db.sh old mode 100644 new mode 100755 diff --git a/Directory.Packages.props b/Directory.Packages.props index cfd8c02161..0bf38f6a86 100644 --- a/Directory.Packages.props +++ b/Directory.Packages.props @@ -36,6 +36,7 @@ + diff --git a/src/Npgsql.OpenTelemetry/Npgsql.OpenTelemetry.csproj b/src/Npgsql.OpenTelemetry/Npgsql.OpenTelemetry.csproj index 7f9fea3eea..543a955db7 100644 --- a/src/Npgsql.OpenTelemetry/Npgsql.OpenTelemetry.csproj +++ b/src/Npgsql.OpenTelemetry/Npgsql.OpenTelemetry.csproj @@ -17,6 +17,10 @@ - + + + + + diff --git a/src/Npgsql.OpenTelemetry/NpgsqlTracingInstrumentation.cs b/src/Npgsql.OpenTelemetry/NpgsqlTracingInstrumentation.cs new file mode 100644 index 0000000000..9f7c047b7f --- /dev/null +++ b/src/Npgsql.OpenTelemetry/NpgsqlTracingInstrumentation.cs @@ -0,0 +1,16 @@ +using System; + +namespace Npgsql.OpenTelemetry; + +sealed class NpgsqlTracingInstrumentation : IDisposable +{ + readonly NpgsqlTracingOptions _originalOptions; + + public NpgsqlTracingInstrumentation(NpgsqlTracingOptions options) + { + _originalOptions = NpgsqlActivitySource.Options; + NpgsqlActivitySource.Options = options; + } + + public void Dispose() => NpgsqlActivitySource.Options = _originalOptions; +} diff --git a/src/Npgsql.OpenTelemetry/TracerProviderBuilderExtensions.cs b/src/Npgsql.OpenTelemetry/TracerProviderBuilderExtensions.cs index 0c34138278..9f5af4f8f3 100644 --- a/src/Npgsql.OpenTelemetry/TracerProviderBuilderExtensions.cs +++ b/src/Npgsql.OpenTelemetry/TracerProviderBuilderExtensions.cs @@ -1,4 +1,5 @@ using System; +using Npgsql.OpenTelemetry; using OpenTelemetry.Trace; // ReSharper disable once CheckNamespace @@ -14,6 +15,12 @@ public static class TracerProviderBuilderExtensions /// public static TracerProviderBuilder AddNpgsql( this TracerProviderBuilder builder, - Action? options = null) - => builder.AddSource("Npgsql"); + Action? configure = null) + { + var options = new NpgsqlTracingOptions(); + configure?.Invoke(options); + return builder + .AddSource("Npgsql") + .AddInstrumentation(() => new NpgsqlTracingInstrumentation(options)); + } } \ No newline at end of file diff --git a/src/Npgsql/NpgsqlActivitySource.cs b/src/Npgsql/NpgsqlActivitySource.cs index 002cf4a638..b5b9777aa9 100644 --- a/src/Npgsql/NpgsqlActivitySource.cs +++ b/src/Npgsql/NpgsqlActivitySource.cs @@ -20,7 +20,9 @@ static NpgsqlActivitySource() internal static bool IsEnabled => Source.HasListeners(); - internal static Activity? CommandStart(NpgsqlConnector connector, string sql) + internal static NpgsqlTracingOptions Options { get; set; } = new(); + + internal static Activity? CommandStart(NpgsqlConnector connector, NpgsqlCommand command) { var settings = connector.Settings; var activity = Source.StartActivity(settings.Database!, ActivityKind.Client); @@ -31,7 +33,7 @@ static NpgsqlActivitySource() activity.SetTag("db.connection_string", connector.UserFacingConnectionString); activity.SetTag("db.user", settings.Username); activity.SetTag("db.name", settings.Database); - activity.SetTag("db.statement", sql); + activity.SetTag("db.statement", command.CommandText); activity.SetTag("db.connection_id", connector.Id); var endPoint = connector.ConnectedEndPoint; @@ -55,6 +57,8 @@ static NpgsqlActivitySource() throw new ArgumentOutOfRangeException("Invalid endpoint type: " + endPoint.GetType()); } + Options.EnrichCommandExecution?.Invoke(activity, "OnStartActivity", command); + return activity; } @@ -64,13 +68,19 @@ internal static void ReceivedFirstResponse(Activity activity) activity.AddEvent(activityEvent); } - internal static void CommandStop(Activity activity) + internal static void CommandStop(Activity activity, NpgsqlCommand command) { activity.SetTag("otel.status_code", "OK"); + activity.SetEndTime(DateTime.UtcNow); + if (activity.IsAllDataRequested) + { + Options.EnrichCommandExecution?.Invoke(activity, "OnStopActivity", command); + } + activity.Dispose(); } - internal static void SetException(Activity activity, Exception ex, bool escaped = true) + internal static void SetException(Activity activity, NpgsqlCommand command, Exception ex, bool escaped = true) { var tags = new ActivityTagsCollection { @@ -83,6 +93,12 @@ internal static void SetException(Activity activity, Exception ex, bool escaped activity.AddEvent(activityEvent); activity.SetTag("otel.status_code", "ERROR"); activity.SetTag("otel.status_description", ex is PostgresException pgEx ? pgEx.SqlState : ex.Message); + activity.SetEndTime(DateTime.UtcNow); + if (activity.IsAllDataRequested) + { + Options.EnrichCommandExecution?.Invoke(activity, "OnException", (command, ex)); + } + activity.Dispose(); } } \ No newline at end of file diff --git a/src/Npgsql/NpgsqlCommand.cs b/src/Npgsql/NpgsqlCommand.cs index e3d07c3ae9..0df3525df8 100644 --- a/src/Npgsql/NpgsqlCommand.cs +++ b/src/Npgsql/NpgsqlCommand.cs @@ -1609,7 +1609,7 @@ internal void TraceCommandStart(NpgsqlConnector connector) { Debug.Assert(CurrentActivity is null); if (NpgsqlActivitySource.IsEnabled) - CurrentActivity = NpgsqlActivitySource.CommandStart(connector, CommandText); + CurrentActivity = NpgsqlActivitySource.CommandStart(connector, this); } internal void TraceReceivedFirstResponse() @@ -1624,7 +1624,7 @@ internal void TraceCommandStop() { if (CurrentActivity is not null) { - NpgsqlActivitySource.CommandStop(CurrentActivity); + NpgsqlActivitySource.CommandStop(CurrentActivity, this); CurrentActivity = null; } } @@ -1633,7 +1633,7 @@ internal void TraceSetException(Exception e) { if (CurrentActivity is not null) { - NpgsqlActivitySource.SetException(CurrentActivity, e); + NpgsqlActivitySource.SetException(CurrentActivity, this, e); CurrentActivity = null; } } diff --git a/src/Npgsql/NpgsqlTracingOptions.cs b/src/Npgsql/NpgsqlTracingOptions.cs index 4aa61beec6..5e198fccac 100644 --- a/src/Npgsql/NpgsqlTracingOptions.cs +++ b/src/Npgsql/NpgsqlTracingOptions.cs @@ -1,9 +1,18 @@ +using System; +using System.Diagnostics; + namespace Npgsql; /// /// Options to configure Npgsql's support for OpenTelemetry tracing. -/// Currently no options are available. /// public class NpgsqlTracingOptions { + /// + /// Gets or sets an action to enrich a Command Execution Activity. + /// + /// + /// + /// + public Action? EnrichCommandExecution { get; set; } } \ No newline at end of file diff --git a/src/Npgsql/Properties/AssemblyInfo.cs b/src/Npgsql/Properties/AssemblyInfo.cs index e71a69a9dd..3201938593 100644 --- a/src/Npgsql/Properties/AssemblyInfo.cs +++ b/src/Npgsql/Properties/AssemblyInfo.cs @@ -46,3 +46,10 @@ "8078a5df97a62d83c9a2db2d072523a8fc491398254c6b89329b8c1dcef43a1e" + "7aa16153bcea2ae9a471145624826f60d7c8e71cd025b554a0177bd935a78096" + "29f0a7afc778ebb4ad033e1bf512c1a9c6ceea26b077bc46cac93800435e77ee")] + +[assembly: InternalsVisibleTo("Npgsql.OpenTelemetry, PublicKey=" + +"0024000004800000940000000602000000240000525341310004000001000100" + +"2b3c590b2a4e3d347e6878dc0ff4d21eb056a50420250c6617044330701d35c9" + +"8078a5df97a62d83c9a2db2d072523a8fc491398254c6b89329b8c1dcef43a1e" + +"7aa16153bcea2ae9a471145624826f60d7c8e71cd025b554a0177bd935a78096" + +"29f0a7afc778ebb4ad033e1bf512c1a9c6ceea26b077bc46cac93800435e77ee")] diff --git a/src/Npgsql/PublicAPI.Unshipped.txt b/src/Npgsql/PublicAPI.Unshipped.txt index 1d8c5e4c4a..9555b8bf8e 100644 --- a/src/Npgsql/PublicAPI.Unshipped.txt +++ b/src/Npgsql/PublicAPI.Unshipped.txt @@ -15,6 +15,8 @@ Npgsql.NpgsqlDataSourceBuilder.UnmapComposite(string? pgName = null, Npgsql.I Npgsql.NpgsqlDataSourceBuilder.UnmapEnum(string? pgName = null, Npgsql.INpgsqlNameTranslator? nameTranslator = null) -> bool Npgsql.NpgsqlDataSourceBuilder.UsePhysicalConnectionInitializer(System.Action? connectionInitializer, System.Func? connectionInitializerAsync) -> Npgsql.NpgsqlDataSourceBuilder! Npgsql.NpgsqlLoggingConfiguration +Npgsql.NpgsqlTracingOptions.EnrichCommandExecution.get -> System.Action? +Npgsql.NpgsqlTracingOptions.EnrichCommandExecution.set -> void Npgsql.Schema.NpgsqlDbColumn.IsIdentity.get -> bool? Npgsql.Schema.NpgsqlDbColumn.IsIdentity.set -> void Npgsql.StatementType.Call = 11 -> Npgsql.StatementType diff --git a/test/Npgsql.Tests/Npgsql.Tests.csproj b/test/Npgsql.Tests/Npgsql.Tests.csproj index 7952ad6301..c8bb528be8 100644 --- a/test/Npgsql.Tests/Npgsql.Tests.csproj +++ b/test/Npgsql.Tests/Npgsql.Tests.csproj @@ -6,8 +6,10 @@ + + diff --git a/test/Npgsql.Tests/OpenTelemetry/NpgsqlTracingOptionsTests.cs b/test/Npgsql.Tests/OpenTelemetry/NpgsqlTracingOptionsTests.cs new file mode 100644 index 0000000000..0920950e11 --- /dev/null +++ b/test/Npgsql.Tests/OpenTelemetry/NpgsqlTracingOptionsTests.cs @@ -0,0 +1,98 @@ +using System; +using System.Collections.Generic; +using System.Diagnostics; +using NUnit.Framework; +using OpenTelemetry; +using OpenTelemetry.Trace; + +namespace Npgsql.Tests.OpenTelemetry; + +[NonParallelizable] +public class NpgsqlTracingOptionsTests : TestBase +{ + [Test] + public void CommandExecution_start_stop() + { + using (var conn = OpenConnection()) + { + conn.ExecuteScalar("SELECT 1"); + } + + Assert.That(_enrichInvocations, Has.Count.EqualTo(2)); + + var (startActivity, startEventName, startObject) = _enrichInvocations[0]; + Assert.That(startEventName, Is.EqualTo("OnStartActivity")); + Assert.That(startObject, Is.TypeOf().With.Property("CommandText").EqualTo("SELECT 1")); + Assert.That(startActivity.Kind, Is.EqualTo(ActivityKind.Client)); + + var (stopActivity, stopEventName, stopObject) = _enrichInvocations[1]; + Assert.That(stopEventName, Is.EqualTo("OnStopActivity")); + Assert.That(stopObject, Is.SameAs(startObject)); + Assert.That(stopActivity, Is.SameAs(startActivity)); + } + + [Test] + public void CommandExecution_start_exception() + { + var exception = Assert.Throws(() => + { + using var conn = OpenConnection(); + conn.ExecuteScalar("BO SELECTA"); + }); + + Assert.That(_enrichInvocations, Has.Count.EqualTo(2)); + + var (startActivity, startEventName, startObject) = _enrichInvocations[0]; + Assert.That(startEventName, Is.EqualTo("OnStartActivity")); + Assert.That(startObject, Is.TypeOf().With.Property("CommandText").EqualTo("BO SELECTA")); + Assert.That(startActivity.Kind, Is.EqualTo(ActivityKind.Client)); + + var (stopActivity, stopEventName, stopObject) = _enrichInvocations[1]; + Assert.That(stopEventName, Is.EqualTo("OnException")); + Assert.That(stopObject, Is.TypeOf>()); + var (stopCommand, stopException) = (ValueTuple)stopObject; + Assert.That(stopCommand.CommandText, Is.EqualTo("BO SELECTA")); + Assert.That(stopException, Is.SameAs(exception)); + Assert.That(stopActivity, Is.SameAs(startActivity)); + } + + [Test] + public void CommandExecution_start_exception_patternmatch() + { + var exception = Assert.Throws(() => + { + using var conn = OpenConnection(); + conn.ExecuteScalar("BO SELECTA"); + }); + + Assert.That(_enrichInvocations, Has.Count.EqualTo(2)); + var (_, stopEventName, stopObject) = _enrichInvocations[1]; + + switch (stopEventName, stopObject) + { + case ("OnException", (NpgsqlCommand stopCommand, Exception stopException)): + Assert.That(stopCommand.CommandText, Is.EqualTo("BO SELECTA")); + Assert.That(stopException, Is.SameAs(exception)); + break; + default: + Assert.Fail($"{nameof(stopEventName)}: '{stopEventName}', {nameof(stopObject)}.GetType(): '{stopObject.GetType()}'"); + break; + } + } + + [SetUp] + public void SetUp() + { + _enrichInvocations.Clear(); + _tracerProvider = Sdk.CreateTracerProviderBuilder() + .AddNpgsql(o => o.EnrichCommandExecution = (activity, eventName, rawObject) => _enrichInvocations.Add((activity, eventName, rawObject))) + .Build(); + } + + [TearDown] + public void TearDown() => _tracerProvider.Dispose(); + + TracerProvider _tracerProvider = null!; + + readonly List<(Activity activity, string eventName, object rawObject)> _enrichInvocations = new(); +}