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();
+}