From 106ee1ed4948d56254f9c853c67005bfb3022b36 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Martin=20St=C3=BChmer?= Date: Mon, 28 Sep 2026 22:49:08 +0200 Subject: [PATCH 01/14] docs(decisions): propose OpenTelemetry telemetry conventions for Pulse --- ...-28-opentelemetry-telemetry-conventions.md | 79 +++++++++++++++++++ 1 file changed, 79 insertions(+) create mode 100644 decisions/2026-09-28-opentelemetry-telemetry-conventions.md diff --git a/decisions/2026-09-28-opentelemetry-telemetry-conventions.md b/decisions/2026-09-28-opentelemetry-telemetry-conventions.md new file mode 100644 index 00000000..b6c06748 --- /dev/null +++ b/decisions/2026-09-28-opentelemetry-telemetry-conventions.md @@ -0,0 +1,79 @@ +--- +authors: + - Martin Stühmer + +applyTo: + - "src/NetEvolve.Pulse/Interceptors/ActivityAndMetrics*.cs" + - "src/NetEvolve.Pulse/Outbox/OutboxProcessorHostedService.cs" + - "src/NetEvolve.Pulse/Internals/Defaults.cs" + - "src/NetEvolve.Pulse/Internals/OperationInstruments.cs" + +created: 2026-09-28 + +lastModified: 2026-09-28 + +state: proposed + +instructions: | + MUST leave the activity status Unset when a Pulse operation succeeds or a stream consumer stops early; MUST set Error with the exception message only on failure. + MUST set the error.type tag to the exception's full type name on the span, the error counter and the duration histogram of a failed operation, and MUST NOT set it otherwise. + MUST record the stream query duration once for every outcome (completed, faulted, abandoned) and tag the abandoned outcome with pulse.stream.completed=false. + MUST keep the legacy units (ms, requests, events, queries, errors, messages) as default until the opt-in ActivityAndMetricsOptions.UseSemanticConventionUnits becomes the default in a later 0.x release; with the opt-in, durations use seconds (unit s) with explicit bucket advice and counters use UCUM annotations ({request}, {event}, {query}, {error}, {message}). + MUST keep the existing metric and tag names. +--- + +# Decision: OpenTelemetry Telemetry Conventions + +Pulse activities and metrics follow the OpenTelemetry Trace API and semantic conventions for status handling, `error.type` and units. Status and `error.type` changes apply immediately, the unit changes are opt-in during a transition period. + +## Context + +Issue #832 lists the deviations of the built-in telemetry in `NetEvolve.Pulse` from the OpenTelemetry guidance: + +* All three `ActivityAndMetrics*Interceptor` classes set `ActivityStatusCode.Ok` on success. The [Trace API](https://opentelemetry.io/docs/specs/otel/trace/api/#set-status) states: "Generally, Instrumentation Libraries SHOULD NOT set the status code to `Ok`, unless explicitly configured to do so." +* Failed operations carry no `error.type`. [Recording errors](https://opentelemetry.io/docs/specs/semconv/general/recording-errors/) states that spans and metrics SHOULD carry `error.type` on failure and SHOULD NOT include it on success. +* Duration histograms use `ms`, counters use plain words (`requests`, `errors`, `messages`). The [metrics guidelines](https://opentelemetry.io/docs/specs/semconv/general/metrics/#instrument-units) state: "When instruments are measuring durations, seconds (i.e. `s`) SHOULD be used." and ask for UCUM annotations such as `{request}`. + +Issue #831 shows that `ActivityAndMetricsStreamQueryInterceptor` records no duration and no outcome when the consumer stops enumerating early, because the recording code runs after the `try`/`finally` block that holds the `yield return`. + +No earlier decision defines the Pulse telemetry contract. Units are part of the exported metric identity: the OpenTelemetry Prometheus exporter appends the unit to the metric name (for example `_milliseconds` or `_seconds`). Changing a unit therefore breaks existing dashboards and alerts. + +## Decision + +* MUST leave the activity status `Unset` when an operation succeeds and when a stream consumer stops early. MUST set `Error` with the exception message on failure. +* MUST set `error.type` to `Exception.GetType().FullName` on the span, the error counter and the duration histogram of a failed operation. The value MUST be identical on all three. The existing `pulse.exception.*` tags stay. +* MUST record `pulse.stream_query.duration` exactly once per stream query, from the `finally` block of the iterator. A stream the consumer abandons carries `pulse.stream.completed=false` on the histogram and the activity, keeps the status `Unset` and carries no `pulse.success` tag. Completed and faulted streams keep their `pulse.success` semantics. +* MUST treat an exception thrown synchronously by the stream handler delegate like an exception thrown during enumeration. +* MUST offer the new units through `ActivityAndMetricsOptions.UseSemanticConventionUnits` (default `false`). The option applies to the interceptors and to `OutboxProcessorHostedService`: + + | Instrument | Legacy unit | Opt-in unit | + | --- | --- | --- | + | `pulse.requests.total` | `requests` | `{request}` | + | `pulse.events.total` | `events` | `{event}` | + | `pulse.stream_query.total` | `queries` | `{query}` | + | `pulse.request.errors`, `pulse.event.errors`, `pulse.stream_query.errors` | `errors` | `{error}` | + | `pulse.outbox.processed.total`, `pulse.outbox.failed.total`, `pulse.outbox.deadletter.total`, `pulse.outbox.pending` | `messages` | `{message}` | + | `pulse.request.duration`, `pulse.event.duration`, `pulse.stream_query.duration`, `pulse.outbox.processing.duration` | `ms` | `s` | + +* MUST provide explicit histogram bucket boundaries for the `s` unit on .NET 9 and later (`InstrumentAdvice`), using the boundaries of the HTTP semantic conventions: 0.005, 0.01, 0.025, 0.05, 0.075, 0.1, 0.25, 0.5, 0.75, 1, 2.5, 5, 7.5, 10. +* MUST keep the metric names and tag names unchanged. +* Transition plan: a later `0.x` release makes `UseSemanticConventionUnits` the default. The legacy units and the option are removed before `1.0.0`. + +## Consequences + +* Backends treat successful Pulse spans as `Unset` and group failures by `error.type`. +* Abandoned streams appear in the duration histogram, so `pulse.stream_query.total` and `pulse.stream_query.duration` stay consistent. +* Dashboards that filter on `Ok` status must filter on "not `Error`" instead. +* Operators opt into the new units once their dashboards and alerts are migrated. Until then the exported metric names stay as they are. +* A process that mixes both unit modes creates two instruments with the same name on the `NetEvolve.Pulse` meter. The option is meant to be set once per application. + +## Alternatives Considered + +* **Switch the units immediately.** Rejected: issue #832 requires an opt-in during the transition, because the unit is part of the exported metric name. +* **AppContext switch instead of options.** Rejected: a process-wide switch cannot be tested per test, and the repository already configures interceptors through options. +* **Rename the `.total` counters.** Not part of this decision. Renaming breaks every dashboard without an opt-in path and can be decided separately. + +## Related Decisions + +* [Extensibility Interface Evolution Before 1.0](./2026-09-24-extensibility-interface-evolution-pre-1-0.md) - Governs how the changed telemetry contract is announced while the major version is `0`. +* [DateTimeOffset and TimeProvider Usage](./2026-01-21-datetimeoffset-and-timeprovider-usage.md) - Durations are measured with `TimeProvider`. From 59413a57ce527e8ba4d392c0efd17dd577413f27 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Martin=20St=C3=BChmer?= Date: Mon, 28 Sep 2026 23:56:46 +0200 Subject: [PATCH 02/14] test(interceptors): cover early stream stop, unset status and error.type in telemetry interceptors Adds failing tests for an abandoned stream query (duration and outcome not recorded), a stream handler that throws synchronously, the Ok status on success and the missing error.type on spans and metrics. --- .../Pipeline/RequestInterceptorsTests.cs | 8 +- ...ActivityAndMetricsEventInterceptorTests.cs | 107 +++++++- ...tivityAndMetricsRequestInterceptorTests.cs | 104 ++++++- ...tyAndMetricsStreamQueryInterceptorTests.cs | 257 +++++++++++++++++- .../Interceptors/PulseMeasurementCollector.cs | 65 +++++ 5 files changed, 529 insertions(+), 12 deletions(-) create mode 100644 tests/NetEvolve.Pulse.Tests.Unit/Interceptors/PulseMeasurementCollector.cs diff --git a/tests/NetEvolve.Pulse.Tests.Integration/Pipeline/RequestInterceptorsTests.cs b/tests/NetEvolve.Pulse.Tests.Integration/Pipeline/RequestInterceptorsTests.cs index 679827fb..adfe617e 100644 --- a/tests/NetEvolve.Pulse.Tests.Integration/Pipeline/RequestInterceptorsTests.cs +++ b/tests/NetEvolve.Pulse.Tests.Integration/Pipeline/RequestInterceptorsTests.cs @@ -104,7 +104,7 @@ public async Task ActivityAndMetrics_Command_RecordsActivityAndMetrics(Cancellat { _ = await Assert.That(result).IsEqualTo("handled"); _ = await Assert.That(capturedActivity).IsNotNull(); - _ = await Assert.That(capturedActivity!.Status).IsEqualTo(ActivityStatusCode.Ok); + _ = await Assert.That(capturedActivity!.Status).IsEqualTo(ActivityStatusCode.Unset); _ = await Assert.That(capturedActivity.GetTagItem("pulse.request.type")).IsEqualTo("Command"); _ = await Assert.That(Interlocked.Read(ref counterValue)).IsGreaterThanOrEqualTo(1L); } @@ -158,7 +158,7 @@ public async Task ActivityAndMetrics_Query_RecordsActivityAndMetrics(Cancellatio { _ = await Assert.That(result).IsEqualTo(42); _ = await Assert.That(capturedActivity).IsNotNull(); - _ = await Assert.That(capturedActivity!.Status).IsEqualTo(ActivityStatusCode.Ok); + _ = await Assert.That(capturedActivity!.Status).IsEqualTo(ActivityStatusCode.Unset); _ = await Assert.That(capturedActivity.GetTagItem("pulse.request.type")).IsEqualTo("Query"); } } @@ -228,7 +228,7 @@ await mediator using (Assert.Multiple()) { _ = await Assert.That(capturedActivity).IsNotNull(); - _ = await Assert.That(capturedActivity!.Status).IsEqualTo(ActivityStatusCode.Ok); + _ = await Assert.That(capturedActivity!.Status).IsEqualTo(ActivityStatusCode.Unset); _ = await Assert.That(Interlocked.Read(ref counterValue)).IsGreaterThanOrEqualTo(1L); } } @@ -291,7 +291,7 @@ var item in mediator { _ = await Assert.That(items).IsEquivalentTo([1, 2, 3]); _ = await Assert.That(capturedActivity).IsNotNull(); - _ = await Assert.That(capturedActivity!.Status).IsEqualTo(ActivityStatusCode.Ok); + _ = await Assert.That(capturedActivity!.Status).IsEqualTo(ActivityStatusCode.Unset); } } diff --git a/tests/NetEvolve.Pulse.Tests.Unit/Interceptors/ActivityAndMetricsEventInterceptorTests.cs b/tests/NetEvolve.Pulse.Tests.Unit/Interceptors/ActivityAndMetricsEventInterceptorTests.cs index d8e765a9..ea8af9d2 100644 --- a/tests/NetEvolve.Pulse.Tests.Unit/Interceptors/ActivityAndMetricsEventInterceptorTests.cs +++ b/tests/NetEvolve.Pulse.Tests.Unit/Interceptors/ActivityAndMetricsEventInterceptorTests.cs @@ -54,7 +54,7 @@ await interceptor [Test] [NotInParallel] - public async Task HandleAsync_WhenHandlerSucceeds_SetsActivityStatusToOk(CancellationToken cancellationToken) + public async Task HandleAsync_WhenHandlerSucceeds_LeavesActivityStatusUnset(CancellationToken cancellationToken) { cancellationToken.ThrowIfCancellationRequested(); @@ -77,7 +77,7 @@ public async Task HandleAsync_WhenHandlerSucceeds_SetsActivityStatusToOk(Cancell using (Assert.Multiple()) { _ = await Assert.That(capturedActivity).IsNotNull(); - _ = await Assert.That(capturedActivity!.Status).IsEqualTo(ActivityStatusCode.Ok); + _ = await Assert.That(capturedActivity!.Status).IsEqualTo(ActivityStatusCode.Unset); #pragma warning disable CS8605 // Unboxing a possibly null value. _ = await Assert.That((bool)capturedActivity.GetTagItem("pulse.success")).IsTrue(); #pragma warning restore CS8605 // Unboxing a possibly null value. @@ -276,6 +276,109 @@ public async Task HandleAsync_WithNonNullCausationId_TagsCausationId(Cancellatio } } + [Test] + [NotInParallel] + public async Task HandleAsync_WhenHandlerThrows_SetsErrorTypeOnActivityAndMetrics( + CancellationToken cancellationToken + ) + { + cancellationToken.ThrowIfCancellationRequested(); + + using var listener = new ActivityListener + { + ShouldListenTo = source => string.Equals(source.Name, "NetEvolve.Pulse", StringComparison.Ordinal), + Sample = (ref _) => ActivitySamplingResult.AllDataAndRecorded, + }; + ActivitySource.AddActivityListener(listener); + using var collector = new PulseMeasurementCollector("pulse.event.name", nameof(MeasuredEvent)); + + var interceptor = new ActivityAndMetricsEventInterceptor(TimeProvider.System); + Activity? capturedActivity = null; + + listener.ActivityStopped = activity => + { + if (string.Equals(activity.DisplayName, "Event.MeasuredEvent", StringComparison.Ordinal)) + { + capturedActivity = activity; + } + }; + + _ = await Assert.ThrowsAsync(async () => + await interceptor + .HandleAsync( + new MeasuredEvent(), + (_, _) => throw new InvalidOperationException("boom"), + cancellationToken + ) + .ConfigureAwait(false) + ); + + var errors = collector.For("pulse.event.errors"); + var durations = collector.For("pulse.event.duration"); + + using (Assert.Multiple()) + { + _ = await Assert.That(capturedActivity).IsNotNull(); + _ = await Assert + .That(capturedActivity!.GetTagItem("error.type")) + .IsEqualTo("System.InvalidOperationException"); + _ = await Assert.That(errors).Count().IsEqualTo(1); + _ = await Assert.That(errors[0].Tags["error.type"]).IsEqualTo("System.InvalidOperationException"); + _ = await Assert.That(durations).Count().IsEqualTo(1); + _ = await Assert.That(durations[0].Tags["error.type"]).IsEqualTo("System.InvalidOperationException"); + } + } + + [Test] + [NotInParallel] + public async Task HandleAsync_WhenHandlerSucceeds_DoesNotSetErrorType(CancellationToken cancellationToken) + { + cancellationToken.ThrowIfCancellationRequested(); + + using var listener = new ActivityListener + { + ShouldListenTo = source => string.Equals(source.Name, "NetEvolve.Pulse", StringComparison.Ordinal), + Sample = (ref _) => ActivitySamplingResult.AllDataAndRecorded, + }; + ActivitySource.AddActivityListener(listener); + using var collector = new PulseMeasurementCollector("pulse.event.name", nameof(MeasuredEvent)); + + var interceptor = new ActivityAndMetricsEventInterceptor(TimeProvider.System); + Activity? capturedActivity = null; + + listener.ActivityStopped = activity => + { + if (string.Equals(activity.DisplayName, "Event.MeasuredEvent", StringComparison.Ordinal)) + { + capturedActivity = activity; + } + }; + + await interceptor + .HandleAsync(new MeasuredEvent(), (_, _) => Task.CompletedTask, cancellationToken) + .ConfigureAwait(false); + + var durations = collector.For("pulse.event.duration"); + + using (Assert.Multiple()) + { + _ = await Assert.That(capturedActivity).IsNotNull(); + _ = await Assert.That(capturedActivity!.GetTagItem("error.type")).IsNull(); + _ = await Assert.That(durations).Count().IsEqualTo(1); + _ = await Assert.That(durations[0].Tags.ContainsKey("error.type")).IsFalse(); + _ = await Assert.That(collector.For("pulse.event.errors")).IsEmpty(); + } + } + + private sealed class MeasuredEvent : IEvent + { + public string Id { get; init; } = Guid.NewGuid().ToString(); + public string? CausationId { get; set; } + public string? CorrelationId { get; set; } + + DateTimeOffset? IEvent.PublishedAt { get; set; } + } + private sealed class TestEvent : IEvent { public string Id { get; init; } = Guid.NewGuid().ToString(); diff --git a/tests/NetEvolve.Pulse.Tests.Unit/Interceptors/ActivityAndMetricsRequestInterceptorTests.cs b/tests/NetEvolve.Pulse.Tests.Unit/Interceptors/ActivityAndMetricsRequestInterceptorTests.cs index 973ae329..55d40d6d 100644 --- a/tests/NetEvolve.Pulse.Tests.Unit/Interceptors/ActivityAndMetricsRequestInterceptorTests.cs +++ b/tests/NetEvolve.Pulse.Tests.Unit/Interceptors/ActivityAndMetricsRequestInterceptorTests.cs @@ -126,7 +126,7 @@ public async Task HandleAsync_WithGenericRequest_CreatesActivityWithCorrectTags( [Test] [NotInParallel] - public async Task HandleAsync_WhenHandlerSucceeds_SetsActivityStatusToOk(CancellationToken cancellationToken) + public async Task HandleAsync_WhenHandlerSucceeds_LeavesActivityStatusUnset(CancellationToken cancellationToken) { cancellationToken.ThrowIfCancellationRequested(); @@ -151,7 +151,7 @@ public async Task HandleAsync_WhenHandlerSucceeds_SetsActivityStatusToOk(Cancell using (Assert.Multiple()) { _ = await Assert.That(capturedActivity).IsNotNull(); - _ = await Assert.That(capturedActivity!.Status).IsEqualTo(ActivityStatusCode.Ok); + _ = await Assert.That(capturedActivity!.Status).IsEqualTo(ActivityStatusCode.Unset); #pragma warning disable CS8605 // Unboxing a possibly null value. _ = await Assert.That((bool)capturedActivity.GetTagItem("pulse.success")).IsTrue(); #pragma warning restore CS8605 // Unboxing a possibly null value. @@ -321,6 +321,106 @@ public async Task HandleAsync_WithNonNullCausationId_TagsCausationId(Cancellatio } } + [Test] + [NotInParallel] + public async Task HandleAsync_WhenHandlerThrows_SetsErrorTypeOnActivityAndMetrics( + CancellationToken cancellationToken + ) + { + cancellationToken.ThrowIfCancellationRequested(); + + using var listener = new ActivityListener + { + ShouldListenTo = source => string.Equals(source.Name, "NetEvolve.Pulse", StringComparison.Ordinal), + Sample = (ref _) => ActivitySamplingResult.AllDataAndRecorded, + }; + ActivitySource.AddActivityListener(listener); + using var collector = new PulseMeasurementCollector("pulse.request.name", nameof(MeasuredCommand)); + + var interceptor = new ActivityAndMetricsRequestInterceptor(TimeProvider.System); + Activity? capturedActivity = null; + + listener.ActivityStopped = activity => + { + if (string.Equals(activity.DisplayName, "Command.MeasuredCommand", StringComparison.Ordinal)) + { + capturedActivity = activity; + } + }; + + _ = await Assert.ThrowsAsync(async () => + await interceptor + .HandleAsync( + new MeasuredCommand(), + (_, _) => throw new InvalidOperationException("boom"), + cancellationToken + ) + .ConfigureAwait(false) + ); + + var errors = collector.For("pulse.request.errors"); + var durations = collector.For("pulse.request.duration"); + + using (Assert.Multiple()) + { + _ = await Assert.That(capturedActivity).IsNotNull(); + _ = await Assert + .That(capturedActivity!.GetTagItem("error.type")) + .IsEqualTo("System.InvalidOperationException"); + _ = await Assert.That(errors).Count().IsEqualTo(1); + _ = await Assert.That(errors[0].Tags["error.type"]).IsEqualTo("System.InvalidOperationException"); + _ = await Assert.That(durations).Count().IsEqualTo(1); + _ = await Assert.That(durations[0].Tags["error.type"]).IsEqualTo("System.InvalidOperationException"); + } + } + + [Test] + [NotInParallel] + public async Task HandleAsync_WhenHandlerSucceeds_DoesNotSetErrorType(CancellationToken cancellationToken) + { + cancellationToken.ThrowIfCancellationRequested(); + + using var listener = new ActivityListener + { + ShouldListenTo = source => string.Equals(source.Name, "NetEvolve.Pulse", StringComparison.Ordinal), + Sample = (ref _) => ActivitySamplingResult.AllDataAndRecorded, + }; + ActivitySource.AddActivityListener(listener); + using var collector = new PulseMeasurementCollector("pulse.request.name", nameof(MeasuredCommand)); + + var interceptor = new ActivityAndMetricsRequestInterceptor(TimeProvider.System); + Activity? capturedActivity = null; + + listener.ActivityStopped = activity => + { + if (string.Equals(activity.DisplayName, "Command.MeasuredCommand", StringComparison.Ordinal)) + { + capturedActivity = activity; + } + }; + + _ = await interceptor + .HandleAsync(new MeasuredCommand(), (_, _) => Task.FromResult("ok"), cancellationToken) + .ConfigureAwait(false); + + var durations = collector.For("pulse.request.duration"); + + using (Assert.Multiple()) + { + _ = await Assert.That(capturedActivity).IsNotNull(); + _ = await Assert.That(capturedActivity!.GetTagItem("error.type")).IsNull(); + _ = await Assert.That(durations).Count().IsEqualTo(1); + _ = await Assert.That(durations[0].Tags.ContainsKey("error.type")).IsFalse(); + _ = await Assert.That(collector.For("pulse.request.errors")).IsEmpty(); + } + } + + private sealed class MeasuredCommand : ICommand + { + public string? CausationId { get; set; } + public string? CorrelationId { get; set; } + } + private sealed class TestCommand : ICommand { public string? CausationId { get; set; } diff --git a/tests/NetEvolve.Pulse.Tests.Unit/Interceptors/ActivityAndMetricsStreamQueryInterceptorTests.cs b/tests/NetEvolve.Pulse.Tests.Unit/Interceptors/ActivityAndMetricsStreamQueryInterceptorTests.cs index 6ec09d89..3abcc0a6 100644 --- a/tests/NetEvolve.Pulse.Tests.Unit/Interceptors/ActivityAndMetricsStreamQueryInterceptorTests.cs +++ b/tests/NetEvolve.Pulse.Tests.Unit/Interceptors/ActivityAndMetricsStreamQueryInterceptorTests.cs @@ -51,7 +51,7 @@ var _ in interceptor [Test] [NotInParallel] - public async Task HandleAsync_WhenStreamCompletes_SetsActivityStatusToOk(CancellationToken cancellationToken) + public async Task HandleAsync_WhenStreamCompletes_LeavesActivityStatusUnset(CancellationToken cancellationToken) { cancellationToken.ThrowIfCancellationRequested(); @@ -81,7 +81,7 @@ var _ in interceptor using (Assert.Multiple()) { _ = await Assert.That(capturedActivity).IsNotNull(); - _ = await Assert.That(capturedActivity!.Status).IsEqualTo(ActivityStatusCode.Ok); + _ = await Assert.That(capturedActivity!.Status).IsEqualTo(ActivityStatusCode.Unset); #pragma warning disable CS8605 // Unboxing a possibly null value. _ = await Assert.That((bool)capturedActivity.GetTagItem("pulse.success")).IsTrue(); #pragma warning restore CS8605 // Unboxing a possibly null value. @@ -142,7 +142,7 @@ var _ in interceptor [Test] [NotInParallel] - public async Task HandleAsync_WithEmptyStream_SetsActivityStatusToOk(CancellationToken cancellationToken) + public async Task HandleAsync_WithEmptyStream_LeavesActivityStatusUnset(CancellationToken cancellationToken) { cancellationToken.ThrowIfCancellationRequested(); @@ -172,7 +172,7 @@ var _ in interceptor using (Assert.Multiple()) { _ = await Assert.That(capturedActivity).IsNotNull(); - _ = await Assert.That(capturedActivity!.Status).IsEqualTo(ActivityStatusCode.Ok); + _ = await Assert.That(capturedActivity!.Status).IsEqualTo(ActivityStatusCode.Unset); #pragma warning disable CS8605 // Unboxing a possibly null value. _ = await Assert.That((bool)capturedActivity.GetTagItem("pulse.success")).IsTrue(); #pragma warning restore CS8605 // Unboxing a possibly null value. @@ -339,6 +339,249 @@ var _ in interceptor } } + [Test] + [NotInParallel("PulseStreamQueryMetrics")] + public async Task HandleAsync_WhenConsumerStopsEarly_RecordsDurationOnceWithAbandonedOutcome( + CancellationToken cancellationToken + ) + { + cancellationToken.ThrowIfCancellationRequested(); + + using var collector = new PulseMeasurementCollector("pulse.request.name", nameof(MeasuredStreamQuery)); + var interceptor = new ActivityAndMetricsStreamQueryInterceptor(TimeProvider.System); + + var consumed = 0; + await foreach ( + var _ in interceptor + .HandleAsync(new MeasuredStreamQuery(), (_, ct) => Items([1, 2, 3], ct), cancellationToken) + .ConfigureAwait(false) + ) + { + consumed++; + if (consumed == 1) + { + break; + } + } + + var durations = collector.For("pulse.stream_query.duration"); + + using (Assert.Multiple()) + { + _ = await Assert.That(durations).Count().IsEqualTo(1); + _ = await Assert.That(durations[0].Tags["pulse.stream.completed"] is false).IsTrue(); + _ = await Assert.That(durations[0].Tags.ContainsKey("pulse.success")).IsFalse(); + _ = await Assert.That(durations[0].Tags.ContainsKey("error.type")).IsFalse(); + _ = await Assert.That(collector.For("pulse.stream_query.errors")).IsEmpty(); + } + } + + [Test] + [NotInParallel] + public async Task HandleAsync_WhenConsumerStopsEarly_TagsActivityAbandonedAndLeavesStatusUnset( + CancellationToken cancellationToken + ) + { + cancellationToken.ThrowIfCancellationRequested(); + + using var listener = new ActivityListener + { + ShouldListenTo = source => string.Equals(source.Name, "NetEvolve.Pulse", StringComparison.Ordinal), + Sample = (ref _) => ActivitySamplingResult.AllDataAndRecorded, + }; + ActivitySource.AddActivityListener(listener); + + var interceptor = new ActivityAndMetricsStreamQueryInterceptor(TimeProvider.System); + Activity? capturedActivity = null; + + listener.ActivityStopped = activity => + { + if (string.Equals(activity.DisplayName, "StreamQuery.MeasuredStreamQuery", StringComparison.Ordinal)) + { + capturedActivity = activity; + } + }; + + var consumed = 0; + await foreach ( + var _ in interceptor + .HandleAsync(new MeasuredStreamQuery(), (_, ct) => Items([1, 2, 3], ct), cancellationToken) + .ConfigureAwait(false) + ) + { + consumed++; + if (consumed == 1) + { + break; + } + } + + using (Assert.Multiple()) + { + _ = await Assert.That(capturedActivity).IsNotNull(); + _ = await Assert.That(capturedActivity!.Status).IsEqualTo(ActivityStatusCode.Unset); + _ = await Assert.That(capturedActivity.GetTagItem("pulse.stream.completed") is false).IsTrue(); + _ = await Assert.That(capturedActivity.GetTagItem("pulse.success")).IsNull(); + _ = await Assert.That(capturedActivity.GetTagItem("error.type")).IsNull(); + } + } + + [Test] + [NotInParallel("PulseStreamQueryMetrics")] + public async Task HandleAsync_WhenStreamCompletes_RecordsDurationOnceWithSuccess( + CancellationToken cancellationToken + ) + { + cancellationToken.ThrowIfCancellationRequested(); + + using var collector = new PulseMeasurementCollector("pulse.request.name", nameof(MeasuredStreamQuery)); + var interceptor = new ActivityAndMetricsStreamQueryInterceptor(TimeProvider.System); + + await foreach ( + var _ in interceptor + .HandleAsync(new MeasuredStreamQuery(), (_, ct) => Items([1, 2, 3], ct), cancellationToken) + .ConfigureAwait(false) + ) + { + // consume items + } + + var durations = collector.For("pulse.stream_query.duration"); + + using (Assert.Multiple()) + { + _ = await Assert.That(durations).Count().IsEqualTo(1); + _ = await Assert.That(durations[0].Tags["pulse.success"] is true).IsTrue(); + _ = await Assert.That(durations[0].Tags.ContainsKey("pulse.stream.completed")).IsFalse(); + _ = await Assert.That(durations[0].Tags.ContainsKey("error.type")).IsFalse(); + _ = await Assert.That(collector.For("pulse.stream_query.total")).Count().IsEqualTo(1); + _ = await Assert.That(collector.For("pulse.stream_query.errors")).IsEmpty(); + } + } + + [Test] + [NotInParallel("PulseStreamQueryMetrics")] + public async Task HandleAsync_WhenStreamFaults_RecordsErrorTypeOnMetrics(CancellationToken cancellationToken) + { + cancellationToken.ThrowIfCancellationRequested(); + + using var collector = new PulseMeasurementCollector("pulse.request.name", nameof(MeasuredStreamQuery)); + var interceptor = new ActivityAndMetricsStreamQueryInterceptor(TimeProvider.System); + + _ = await Assert.ThrowsAsync(async () => + { + await foreach ( + var _ in interceptor + .HandleAsync( + new MeasuredStreamQuery(), + (_, ct) => ThrowingItems(new InvalidOperationException("boom"), ct), + cancellationToken + ) + .ConfigureAwait(false) + ) + { + // consume items until exception + } + }); + + var durations = collector.For("pulse.stream_query.duration"); + var errors = collector.For("pulse.stream_query.errors"); + + using (Assert.Multiple()) + { + _ = await Assert.That(durations).Count().IsEqualTo(1); + _ = await Assert.That(durations[0].Tags["pulse.success"] is false).IsTrue(); + _ = await Assert.That(durations[0].Tags["error.type"]).IsEqualTo("System.InvalidOperationException"); + _ = await Assert.That(errors).Count().IsEqualTo(1); + _ = await Assert.That(errors[0].Tags["error.type"]).IsEqualTo("System.InvalidOperationException"); + } + } + + [Test] + [NotInParallel("PulseStreamQueryMetrics")] + public async Task HandleAsync_WhenHandlerThrowsSynchronously_RecordsErrorAndDuration( + CancellationToken cancellationToken + ) + { + cancellationToken.ThrowIfCancellationRequested(); + + using var collector = new PulseMeasurementCollector("pulse.request.name", nameof(MeasuredStreamQuery)); + var interceptor = new ActivityAndMetricsStreamQueryInterceptor(TimeProvider.System); + var testException = new InvalidOperationException("sync"); + + var exception = await Assert.ThrowsAsync(async () => + { + await foreach ( + var _ in interceptor + .HandleAsync(new MeasuredStreamQuery(), (_, _) => throw testException, cancellationToken) + .ConfigureAwait(false) + ) + { + // no items expected + } + }); + + var durations = collector.For("pulse.stream_query.duration"); + + using (Assert.Multiple()) + { + _ = await Assert.That(exception).IsSameReferenceAs(testException); + _ = await Assert.That(collector.For("pulse.stream_query.errors")).Count().IsEqualTo(1); + _ = await Assert.That(durations).Count().IsEqualTo(1); + _ = await Assert.That(durations[0].Tags["pulse.success"] is false).IsTrue(); + } + } + + [Test] + [NotInParallel] + public async Task HandleAsync_WhenHandlerThrows_SetsErrorTypeOnActivity(CancellationToken cancellationToken) + { + cancellationToken.ThrowIfCancellationRequested(); + + using var listener = new ActivityListener + { + ShouldListenTo = source => string.Equals(source.Name, "NetEvolve.Pulse", StringComparison.Ordinal), + Sample = (ref _) => ActivitySamplingResult.AllDataAndRecorded, + }; + ActivitySource.AddActivityListener(listener); + + var interceptor = new ActivityAndMetricsStreamQueryInterceptor(TimeProvider.System); + Activity? capturedActivity = null; + + listener.ActivityStopped = activity => + { + if (string.Equals(activity.DisplayName, "StreamQuery.MeasuredStreamQuery", StringComparison.Ordinal)) + { + capturedActivity = activity; + } + }; + + _ = await Assert.ThrowsAsync(async () => + { + await foreach ( + var _ in interceptor + .HandleAsync( + new MeasuredStreamQuery(), + (_, ct) => ThrowingItems(new InvalidOperationException("boom"), ct), + cancellationToken + ) + .ConfigureAwait(false) + ) + { + // consume items until exception + } + }); + + using (Assert.Multiple()) + { + _ = await Assert.That(capturedActivity).IsNotNull(); + _ = await Assert.That(capturedActivity!.Status).IsEqualTo(ActivityStatusCode.Error); + _ = await Assert + .That(capturedActivity.GetTagItem("error.type")) + .IsEqualTo("System.InvalidOperationException"); + } + } + private static IAsyncEnumerable Items(IEnumerable items, CancellationToken cancellationToken = default) => ItemsCore(items, cancellationToken); @@ -371,4 +614,10 @@ private sealed class TestStreamQuery : IStreamQuery public string? CausationId { get; set; } public string? CorrelationId { get; set; } } + + private sealed class MeasuredStreamQuery : IStreamQuery + { + public string? CausationId { get; set; } + public string? CorrelationId { get; set; } + } } diff --git a/tests/NetEvolve.Pulse.Tests.Unit/Interceptors/PulseMeasurementCollector.cs b/tests/NetEvolve.Pulse.Tests.Unit/Interceptors/PulseMeasurementCollector.cs new file mode 100644 index 00000000..94793410 --- /dev/null +++ b/tests/NetEvolve.Pulse.Tests.Unit/Interceptors/PulseMeasurementCollector.cs @@ -0,0 +1,65 @@ +namespace NetEvolve.Pulse.Tests.Unit.Interceptors; + +using System; +using System.Collections.Concurrent; +using System.Collections.Generic; +using System.Diagnostics.Metrics; +using System.Linq; + +/// +/// Collects measurements of the NetEvolve.Pulse meter that carry a specific tag value, so tests +/// running in parallel with other Pulse telemetry only see their own measurements. +/// +internal sealed class PulseMeasurementCollector : IDisposable +{ + private readonly MeterListener _listener = new(); + private readonly ConcurrentQueue _measurements = new(); + private readonly string _tagKey; + private readonly string _tagValue; + + public PulseMeasurementCollector(string tagKey, string tagValue) + { + _tagKey = tagKey; + _tagValue = tagValue; + + _listener.InstrumentPublished = (instrument, listener) => + { + if (string.Equals(instrument.Meter.Name, "NetEvolve.Pulse", StringComparison.Ordinal)) + { + listener.EnableMeasurementEvents(instrument); + } + }; + _listener.SetMeasurementEventCallback((instrument, value, tags, _) => Add(instrument, value, tags)); + _listener.SetMeasurementEventCallback((instrument, value, tags, _) => Add(instrument, value, tags)); + _listener.Start(); + } + + public IReadOnlyList For(string instrumentName) => + [.. _measurements.Where(m => string.Equals(m.Instrument, instrumentName, StringComparison.Ordinal))]; + + public void Dispose() => _listener.Dispose(); + + private void Add(Instrument instrument, double value, ReadOnlySpan> tags) + { + var tagMap = new Dictionary(StringComparer.Ordinal); + foreach (var tag in tags) + { + tagMap[tag.Key] = tag.Value; + } + + if ( + tagMap.TryGetValue(_tagKey, out var actual) + && string.Equals(actual as string, _tagValue, StringComparison.Ordinal) + ) + { + _measurements.Enqueue(new Measurement(instrument.Name, instrument.Unit, value, tagMap)); + } + } + + internal sealed record Measurement( + string Instrument, + string? Unit, + double Value, + IReadOnlyDictionary Tags + ); +} From c53ce2ce3d867e1c582259c046ed31d70b9d4bd4 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Martin=20St=C3=BChmer?= Date: Tue, 29 Sep 2026 00:38:05 +0200 Subject: [PATCH 03/14] fix(interceptors): record stream query duration and outcome when the consumer stops early Moves the outcome recording of ActivityAndMetricsStreamQueryInterceptor into the finally block of the iterator, so completed, faulted and abandoned streams record pulse.stream_query.duration exactly once. Abandoned streams carry pulse.stream.completed=false and keep the status Unset. A handler that throws synchronously is now recorded as a failure. Successful streams no longer set the Ok status, failures carry error.type, and the unused StreamQueryType tag constant is removed. --- ...treamQueryInterceptor{TQuery,TResponse}.cs | 105 +++++++++++++----- src/NetEvolve.Pulse/Internals/Defaults.cs | 7 +- 2 files changed, 80 insertions(+), 32 deletions(-) diff --git a/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsStreamQueryInterceptor{TQuery,TResponse}.cs b/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsStreamQueryInterceptor{TQuery,TResponse}.cs index d00fba0f..c675065b 100644 --- a/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsStreamQueryInterceptor{TQuery,TResponse}.cs +++ b/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsStreamQueryInterceptor{TQuery,TResponse}.cs @@ -81,6 +81,13 @@ internal sealed class ActivityAndMetricsStreamQueryInterceptorMeasures and records execution duration /// Captures exception details on failure /// Marks success/failure status in both activity and metrics + /// Leaves the activity status unless the stream faults + /// Sets error.type on the activity, the error counter and the duration histogram on failure + /// + /// Records the duration exactly once for every outcome. A stream whose consumer stops enumerating early + /// (for example , Take or a disconnected client) is tagged pulse.stream.completed=false + /// and carries no pulse.success tag + /// /// Yields items unchanged without buffering /// /// @@ -118,13 +125,23 @@ public async IAsyncEnumerable HandleAsync( StreamQueryCounter.Add(1, tags); // yield return is not allowed inside a try/catch block, so we capture any exception - // from the inner enumerator and re-throw it after the yield loop completes. + // from the handler or the inner enumerator and re-throw it after the finally block ran. ExceptionDispatchInfo? caughtExceptionInfo = null; + var completed = false; + IAsyncEnumerator? enumerator = null; - var enumerator = handler(request, cancellationToken).GetAsyncEnumerator(cancellationToken); try { - while (true) + try + { + enumerator = handler(request, cancellationToken).GetAsyncEnumerator(cancellationToken); + } + catch (Exception ex) + { + caughtExceptionInfo = ExceptionDispatchInfo.Capture(ex); + } + + while (enumerator is not null) { bool hasNext; try @@ -139,6 +156,7 @@ public async IAsyncEnumerable HandleAsync( if (!hasNext) { + completed = true; break; } @@ -148,46 +166,73 @@ public async IAsyncEnumerable HandleAsync( } finally { - await enumerator.DisposeAsync().ConfigureAwait(false); + try + { + if (enumerator is not null) + { + await enumerator.DisposeAsync().ConfigureAwait(false); + } + } + finally + { + // Runs for every outcome, including a consumer that stops early and disposes the iterator. + RecordOutcome(activity, tags, startTime, caughtExceptionInfo?.SourceException, completed); + } } - if (caughtExceptionInfo is not null) + caughtExceptionInfo?.Throw(); + } + + /// + /// Records the duration and outcome of a stream query on the activity and the metrics. + /// + /// The activity of the stream query, if sampled. + /// The base tags of the stream query. + /// The time the stream query started. + /// The exception that faulted the stream, or . + /// if the stream was enumerated to its end. + private void RecordOutcome( + Activity? activity, + TagList tags, + DateTimeOffset startTime, + Exception? exception, + bool completed + ) + { + var endTime = _timeProvider.GetUtcNow(); + var duration = (endTime - startTime).TotalMilliseconds; + _ = activity?.SetEndTime(endTime.UtcDateTime); + + if (exception is not null) { - var ex = caughtExceptionInfo.SourceException; - var errorTime = _timeProvider.GetUtcNow(); + var errorType = exception.GetType().FullName; // Capture comprehensive exception details in the activity _ = activity - ?.SetStatus(ActivityStatusCode.Error, ex.Message) - .SetEndTime(errorTime.UtcDateTime) - .SetTag(ExceptionType, ex.GetType().FullName) - .SetTag(ExceptionMessage, ex.Message) - .SetTag(ExceptionStackTrace, ex.StackTrace) - .SetTag(ExceptionTimestamp, errorTime) + ?.SetStatus(ActivityStatusCode.Error, exception.Message) + .SetTag(ErrorType, errorType) + .SetTag(ExceptionType, errorType) + .SetTag(ExceptionMessage, exception.Message) + .SetTag(ExceptionStackTrace, exception.StackTrace) + .SetTag(ExceptionTimestamp, endTime) .SetTag(Success, value: false); - // Increment error counters and record failed execution duration - ErrorsCounter.Add(1, tags); - StreamQueryDurationHistogram.Record( - (errorTime - startTime).TotalMilliseconds, - [.. tags, new(Success, false)] - ); + ErrorsCounter.Add(1, [.. tags, new(ErrorType, errorType)]); + StreamQueryDurationHistogram.Record(duration, [.. tags, new(Success, false), new(ErrorType, errorType)]); + } + else if (completed) + { + // Status stays Unset on success, as required by the OpenTelemetry Trace API. + _ = activity?.SetTag(ResponseTimestamp, endTime).SetTag(Success, value: true); - caughtExceptionInfo.Throw(); + StreamQueryDurationHistogram.Record(duration, [.. tags, new(Success, true)]); } else { - var endTime = _timeProvider.GetUtcNow(); - - // Mark activity as successful - _ = activity - ?.SetStatus(ActivityStatusCode.Ok) - .SetEndTime(endTime.UtcDateTime) - .SetTag(ResponseTimestamp, endTime) - .SetTag(Success, value: true); + // The consumer stopped early: not an error, so the status stays Unset. + _ = activity?.SetTag(StreamCompleted, value: false); - // Record successful execution duration - StreamQueryDurationHistogram.Record((endTime - startTime).TotalMilliseconds, [.. tags, new(Success, true)]); + StreamQueryDurationHistogram.Record(duration, [.. tags, new(StreamCompleted, false)]); } } } diff --git a/src/NetEvolve.Pulse/Internals/Defaults.cs b/src/NetEvolve.Pulse/Internals/Defaults.cs index ea4f09f2..b1103a67 100644 --- a/src/NetEvolve.Pulse/Internals/Defaults.cs +++ b/src/NetEvolve.Pulse/Internals/Defaults.cs @@ -125,10 +125,13 @@ internal static class Tags /// Tag name for stream query causation ID. internal const string StreamQueryCausationId = "pulse.causation_id"; - /// Tag name for stream query type (type name of the query). - internal const string StreamQueryType = "query.type"; + /// Tag name marking a stream query whose consumer stopped enumerating before the stream completed. + internal const string StreamCompleted = "pulse.stream.completed"; // General tags + /// Tag name for the OpenTelemetry error.type attribute, set to the full exception type name on failure. + internal const string ErrorType = "error.type"; + /// Tag name for success/failure indicator. internal const string Success = "pulse.success"; } From 025461daa648a810a9ccb6d91d837699b6cc1d82 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Martin=20St=C3=BChmer?= Date: Tue, 29 Sep 2026 00:38:07 +0200 Subject: [PATCH 04/14] fix(interceptors): leave activity status unset on success and add error.type The request and event telemetry interceptors no longer set ActivityStatusCode.Ok on success and add error.type (the full exception type name) to the span, the error counter and the duration histogram on failure. --- ...tivityAndMetricsEventInterceptor{TEvent}.cs | 18 ++++++++++++------ ...csRequestInterceptor{TRequest,TResponse}.cs | 18 ++++++++++++------ 2 files changed, 24 insertions(+), 12 deletions(-) diff --git a/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsEventInterceptor{TEvent}.cs b/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsEventInterceptor{TEvent}.cs index 3c8a4cc7..b36a4342 100644 --- a/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsEventInterceptor{TEvent}.cs +++ b/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsEventInterceptor{TEvent}.cs @@ -63,6 +63,8 @@ internal sealed class ActivityAndMetricsEventInterceptor : IEventInterce /// Measures and records execution duration /// Captures exception details on failure /// Marks success/failure status in both activity and metrics + /// Leaves the activity status unless the handler fails + /// Sets error.type on the activity, the error counter and the duration histogram on failure /// /// public async Task HandleAsync( @@ -102,10 +104,9 @@ public async Task HandleAsync( var endTime = _timeProvider.GetUtcNow(); - // Mark activity as successful + // Status stays Unset on success, as required by the OpenTelemetry Trace API. _ = activity - ?.SetStatus(ActivityStatusCode.Ok) - .SetEndTime(endTime.UtcDateTime) + ?.SetEndTime(endTime.UtcDateTime) .SetTag(EventCompletionTimestamp, endTime) .SetTag(Success, value: true); @@ -115,20 +116,25 @@ public async Task HandleAsync( catch (Exception ex) { var errorTime = _timeProvider.GetUtcNow(); + var errorType = ex.GetType().FullName; // Capture comprehensive exception details in the activity _ = activity ?.SetStatus(ActivityStatusCode.Error, ex.Message) .SetEndTime(errorTime.UtcDateTime) - .SetTag(ExceptionType, ex.GetType().FullName) + .SetTag(ErrorType, errorType) + .SetTag(ExceptionType, errorType) .SetTag(ExceptionMessage, ex.Message) .SetTag(ExceptionStackTrace, ex.StackTrace) .SetTag(ExceptionTimestamp, errorTime) .SetTag(Success, value: false); // Increment error counters and record failed execution duration - ErrorsCounter.Add(1, tags); - EventDurationHistogram.Record((errorTime - startTime).TotalMilliseconds, [.. tags, new(Success, false)]); + ErrorsCounter.Add(1, [.. tags, new(ErrorType, errorType)]); + EventDurationHistogram.Record( + (errorTime - startTime).TotalMilliseconds, + [.. tags, new(Success, false), new(ErrorType, errorType)] + ); throw; } diff --git a/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsRequestInterceptor{TRequest,TResponse}.cs b/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsRequestInterceptor{TRequest,TResponse}.cs index 942ba238..081bd0a2 100644 --- a/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsRequestInterceptor{TRequest,TResponse}.cs +++ b/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsRequestInterceptor{TRequest,TResponse}.cs @@ -65,6 +65,8 @@ internal sealed class ActivityAndMetricsRequestInterceptor /// Measures and records execution duration /// Captures exception details on failure /// Marks success/failure status in both activity and metrics + /// Leaves the activity status unless the handler fails + /// Sets error.type on the activity, the error counter and the duration histogram on failure /// /// public async Task HandleAsync( @@ -111,10 +113,9 @@ public async Task HandleAsync( var endTime = _timeProvider.GetUtcNow(); - // Mark activity as successful + // Status stays Unset on success, as required by the OpenTelemetry Trace API. _ = activity - ?.SetStatus(ActivityStatusCode.Ok) - .SetEndTime(endTime.UtcDateTime) + ?.SetEndTime(endTime.UtcDateTime) .SetTag(ResponseTimestamp, endTime) .SetTag(Success, value: true); @@ -126,20 +127,25 @@ public async Task HandleAsync( catch (Exception ex) { var errorTime = _timeProvider.GetUtcNow(); + var errorType = ex.GetType().FullName; // Capture comprehensive exception details in the activity _ = activity ?.SetStatus(ActivityStatusCode.Error, ex.Message) .SetEndTime(errorTime.UtcDateTime) - .SetTag(ExceptionType, ex.GetType().FullName) + .SetTag(ErrorType, errorType) + .SetTag(ExceptionType, errorType) .SetTag(ExceptionMessage, ex.Message) .SetTag(ExceptionStackTrace, ex.StackTrace) .SetTag(ExceptionTimestamp, errorTime) .SetTag(Success, value: false); // Increment error counters and record failed execution duration - ErrorsCounter.Add(1, tags); - RequestDurationHistogram.Record((errorTime - startTime).TotalMilliseconds, [.. tags, new(Success, false)]); + ErrorsCounter.Add(1, [.. tags, new(ErrorType, errorType)]); + RequestDurationHistogram.Record( + (errorTime - startTime).TotalMilliseconds, + [.. tags, new(Success, false), new(ErrorType, errorType)] + ); throw; } From 088200c1372dd3d07b5c715a7a8a00cb7815e872 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Martin=20St=C3=BChmer?= Date: Tue, 29 Sep 2026 00:40:18 +0200 Subject: [PATCH 05/14] test(telemetry): specify opt-in OpenTelemetry units for Pulse metrics Adds tests for ActivityAndMetricsOptions.UseSemanticConventionUnits: durations in seconds (unit s) and UCUM annotation units for the interceptor and outbox counters, while the default keeps the legacy units. --- .../ActivityMetricsExtensionsTests.cs | 41 ++++++++ ...ActivityAndMetricsEventInterceptorTests.cs | 80 ++++++++++++++++ ...tivityAndMetricsRequestInterceptorTests.cs | 93 +++++++++++++++++++ ...tyAndMetricsStreamQueryInterceptorTests.cs | 89 ++++++++++++++++++ .../Interceptors/PulseMeasurementCollector.cs | 16 ++-- .../OutboxProcessorHostedServiceTests.cs | 43 +++++++++ 6 files changed, 356 insertions(+), 6 deletions(-) diff --git a/tests/NetEvolve.Pulse.Tests.Unit/ActivityMetricsExtensionsTests.cs b/tests/NetEvolve.Pulse.Tests.Unit/ActivityMetricsExtensionsTests.cs index f2e0040c..8d53bef8 100644 --- a/tests/NetEvolve.Pulse.Tests.Unit/ActivityMetricsExtensionsTests.cs +++ b/tests/NetEvolve.Pulse.Tests.Unit/ActivityMetricsExtensionsTests.cs @@ -4,6 +4,7 @@ using System.Linq; using System.Threading.Tasks; using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Options; using NetEvolve.Extensions.TUnit; using NetEvolve.Pulse.Extensibility; using NetEvolve.Pulse.Interceptors; @@ -109,4 +110,44 @@ public async Task AddActivityAndMetrics_ReturnsSameBuilder() _ = await Assert.That(result).IsTypeOf(); } } + + [Test] + public async Task AddActivityAndMetrics_WithNullConfigure_ThrowsArgumentNullException() + { + var builder = new MediatorBuilder(new ServiceCollection()); + + _ = await Assert.That(() => builder.AddActivityAndMetrics(configure: null!)).Throws(); + } + + [Test] + public async Task AddActivityAndMetrics_WithoutConfigure_DefaultsToLegacyUnits() + { + var services = new ServiceCollection(); + _ = new MediatorBuilder(services).AddActivityAndMetrics(); + + await using var provider = services.BuildServiceProvider(); + + var options = provider.GetRequiredService>().Value; + + _ = await Assert.That(options.UseSemanticConventionUnits).IsFalse(); + } + + [Test] + public async Task AddActivityAndMetrics_WithConfigure_AppliesOptions() + { + var services = new ServiceCollection(); + var builder = new MediatorBuilder(services); + + var result = builder.AddActivityAndMetrics(options => options.UseSemanticConventionUnits = true); + + await using var provider = services.BuildServiceProvider(); + + var options = provider.GetRequiredService>().Value; + + using (Assert.Multiple()) + { + _ = await Assert.That(result).IsSameReferenceAs(builder); + _ = await Assert.That(options.UseSemanticConventionUnits).IsTrue(); + } + } } diff --git a/tests/NetEvolve.Pulse.Tests.Unit/Interceptors/ActivityAndMetricsEventInterceptorTests.cs b/tests/NetEvolve.Pulse.Tests.Unit/Interceptors/ActivityAndMetricsEventInterceptorTests.cs index ea8af9d2..10cc8462 100644 --- a/tests/NetEvolve.Pulse.Tests.Unit/Interceptors/ActivityAndMetricsEventInterceptorTests.cs +++ b/tests/NetEvolve.Pulse.Tests.Unit/Interceptors/ActivityAndMetricsEventInterceptorTests.cs @@ -1,6 +1,8 @@ namespace NetEvolve.Pulse.Tests.Unit.Interceptors; using System.Diagnostics; +using Microsoft.Extensions.Options; +using Microsoft.Extensions.Time.Testing; using NetEvolve.Extensions.TUnit; using NetEvolve.Pulse.Extensibility; using NetEvolve.Pulse.Interceptors; @@ -370,6 +372,84 @@ await interceptor } } + [Test] + [NotInParallel] + public async Task HandleAsync_WithSemanticConventionUnits_RecordsSecondsAndUcumUnits( + CancellationToken cancellationToken + ) + { + cancellationToken.ThrowIfCancellationRequested(); + + using var collector = new PulseMeasurementCollector("pulse.event.name", nameof(MeasuredEvent)); + var timeProvider = new FakeTimeProvider(); + var interceptor = new ActivityAndMetricsEventInterceptor( + timeProvider, + Options.Create(new ActivityAndMetricsOptions { UseSemanticConventionUnits = true }) + ); + + await interceptor + .HandleAsync( + new MeasuredEvent(), + (_, _) => + { + timeProvider.Advance(TimeSpan.FromMilliseconds(250)); + return Task.CompletedTask; + }, + cancellationToken + ) + .ConfigureAwait(false); + _ = await Assert.ThrowsAsync(async () => + await interceptor + .HandleAsync( + new MeasuredEvent(), + (_, _) => throw new InvalidOperationException("boom"), + cancellationToken + ) + .ConfigureAwait(false) + ); + + var durations = collector.For("pulse.event.duration"); + + using (Assert.Multiple()) + { + _ = await Assert.That(durations).Count().IsEqualTo(2); + _ = await Assert.That(durations[0].Unit).IsEqualTo("s"); + _ = await Assert.That(durations[0].Value).IsEqualTo(0.25); + _ = await Assert.That(collector.For("pulse.events.total")[0].Unit).IsEqualTo("{event}"); + _ = await Assert.That(collector.For("pulse.event.errors")[0].Unit).IsEqualTo("{error}"); + } + } + + [Test] + [NotInParallel] + public async Task HandleAsync_WithDefaultOptions_KeepsLegacyUnits(CancellationToken cancellationToken) + { + cancellationToken.ThrowIfCancellationRequested(); + + using var collector = new PulseMeasurementCollector("pulse.event.name", nameof(MeasuredEvent)); + var interceptor = new ActivityAndMetricsEventInterceptor(new FakeTimeProvider()); + + await interceptor + .HandleAsync(new MeasuredEvent(), (_, _) => Task.CompletedTask, cancellationToken) + .ConfigureAwait(false); + _ = await Assert.ThrowsAsync(async () => + await interceptor + .HandleAsync( + new MeasuredEvent(), + (_, _) => throw new InvalidOperationException("boom"), + cancellationToken + ) + .ConfigureAwait(false) + ); + + using (Assert.Multiple()) + { + _ = await Assert.That(collector.For("pulse.event.duration")[0].Unit).IsEqualTo("ms"); + _ = await Assert.That(collector.For("pulse.events.total")[0].Unit).IsEqualTo("events"); + _ = await Assert.That(collector.For("pulse.event.errors")[0].Unit).IsEqualTo("errors"); + } + } + private sealed class MeasuredEvent : IEvent { public string Id { get; init; } = Guid.NewGuid().ToString(); diff --git a/tests/NetEvolve.Pulse.Tests.Unit/Interceptors/ActivityAndMetricsRequestInterceptorTests.cs b/tests/NetEvolve.Pulse.Tests.Unit/Interceptors/ActivityAndMetricsRequestInterceptorTests.cs index 55d40d6d..0f5ef228 100644 --- a/tests/NetEvolve.Pulse.Tests.Unit/Interceptors/ActivityAndMetricsRequestInterceptorTests.cs +++ b/tests/NetEvolve.Pulse.Tests.Unit/Interceptors/ActivityAndMetricsRequestInterceptorTests.cs @@ -1,6 +1,8 @@ namespace NetEvolve.Pulse.Tests.Unit.Interceptors; using System.Diagnostics; +using Microsoft.Extensions.Options; +using Microsoft.Extensions.Time.Testing; using NetEvolve.Extensions.TUnit; using NetEvolve.Pulse.Extensibility; using NetEvolve.Pulse.Interceptors; @@ -415,6 +417,97 @@ public async Task HandleAsync_WhenHandlerSucceeds_DoesNotSetErrorType(Cancellati } } + [Test] + [NotInParallel] + public async Task HandleAsync_WithSemanticConventionUnits_RecordsSecondsAndUcumUnits( + CancellationToken cancellationToken + ) + { + cancellationToken.ThrowIfCancellationRequested(); + + using var collector = new PulseMeasurementCollector("pulse.request.name", nameof(MeasuredCommand)); + var timeProvider = new FakeTimeProvider(); + var interceptor = new ActivityAndMetricsRequestInterceptor( + timeProvider, + Options.Create(new ActivityAndMetricsOptions { UseSemanticConventionUnits = true }) + ); + + _ = await interceptor + .HandleAsync( + new MeasuredCommand(), + (_, _) => + { + timeProvider.Advance(TimeSpan.FromMilliseconds(1500)); + return Task.FromResult("ok"); + }, + cancellationToken + ) + .ConfigureAwait(false); + _ = await Assert.ThrowsAsync(async () => + await interceptor + .HandleAsync( + new MeasuredCommand(), + (_, _) => throw new InvalidOperationException("boom"), + cancellationToken + ) + .ConfigureAwait(false) + ); + + var durations = collector.For("pulse.request.duration"); + + using (Assert.Multiple()) + { + _ = await Assert.That(durations).Count().IsEqualTo(2); + _ = await Assert.That(durations[0].Unit).IsEqualTo("s"); + _ = await Assert.That(durations[0].Value).IsEqualTo(1.5); + _ = await Assert.That(collector.For("pulse.requests.total")[0].Unit).IsEqualTo("{request}"); + _ = await Assert.That(collector.For("pulse.request.errors")[0].Unit).IsEqualTo("{error}"); + } + } + + [Test] + [NotInParallel] + public async Task HandleAsync_WithDefaultOptions_KeepsLegacyUnits(CancellationToken cancellationToken) + { + cancellationToken.ThrowIfCancellationRequested(); + + using var collector = new PulseMeasurementCollector("pulse.request.name", nameof(MeasuredCommand)); + var timeProvider = new FakeTimeProvider(); + var interceptor = new ActivityAndMetricsRequestInterceptor(timeProvider); + + _ = await interceptor + .HandleAsync( + new MeasuredCommand(), + (_, _) => + { + timeProvider.Advance(TimeSpan.FromMilliseconds(1500)); + return Task.FromResult("ok"); + }, + cancellationToken + ) + .ConfigureAwait(false); + _ = await Assert.ThrowsAsync(async () => + await interceptor + .HandleAsync( + new MeasuredCommand(), + (_, _) => throw new InvalidOperationException("boom"), + cancellationToken + ) + .ConfigureAwait(false) + ); + + var durations = collector.For("pulse.request.duration"); + + using (Assert.Multiple()) + { + _ = await Assert.That(durations).Count().IsEqualTo(2); + _ = await Assert.That(durations[0].Unit).IsEqualTo("ms"); + _ = await Assert.That(durations[0].Value).IsEqualTo(1500d); + _ = await Assert.That(collector.For("pulse.requests.total")[0].Unit).IsEqualTo("requests"); + _ = await Assert.That(collector.For("pulse.request.errors")[0].Unit).IsEqualTo("errors"); + } + } + private sealed class MeasuredCommand : ICommand { public string? CausationId { get; set; } diff --git a/tests/NetEvolve.Pulse.Tests.Unit/Interceptors/ActivityAndMetricsStreamQueryInterceptorTests.cs b/tests/NetEvolve.Pulse.Tests.Unit/Interceptors/ActivityAndMetricsStreamQueryInterceptorTests.cs index 3abcc0a6..7475715a 100644 --- a/tests/NetEvolve.Pulse.Tests.Unit/Interceptors/ActivityAndMetricsStreamQueryInterceptorTests.cs +++ b/tests/NetEvolve.Pulse.Tests.Unit/Interceptors/ActivityAndMetricsStreamQueryInterceptorTests.cs @@ -2,6 +2,8 @@ namespace NetEvolve.Pulse.Tests.Unit.Interceptors; using System.Diagnostics; using System.Runtime.CompilerServices; +using Microsoft.Extensions.Options; +using Microsoft.Extensions.Time.Testing; using NetEvolve.Extensions.TUnit; using NetEvolve.Pulse.Extensibility; using NetEvolve.Pulse.Interceptors; @@ -582,6 +584,93 @@ var _ in interceptor } } + [Test] + [NotInParallel("PulseStreamQueryMetrics")] + public async Task HandleAsync_WithSemanticConventionUnits_RecordsSecondsAndUcumUnits( + CancellationToken cancellationToken + ) + { + cancellationToken.ThrowIfCancellationRequested(); + + using var collector = new PulseMeasurementCollector("pulse.request.name", nameof(MeasuredStreamQuery)); + var timeProvider = new FakeTimeProvider(); + var interceptor = new ActivityAndMetricsStreamQueryInterceptor( + timeProvider, + Options.Create(new ActivityAndMetricsOptions { UseSemanticConventionUnits = true }) + ); + + await foreach ( + var _ in interceptor + .HandleAsync(new MeasuredStreamQuery(), (_, ct) => Items([1, 2], ct), cancellationToken) + .ConfigureAwait(false) + ) + { + timeProvider.Advance(TimeSpan.FromSeconds(1)); + } + + _ = await Assert.ThrowsAsync(async () => + { + await foreach ( + var _ in interceptor + .HandleAsync( + new MeasuredStreamQuery(), + (_, ct) => ThrowingItems(new InvalidOperationException("boom"), ct), + cancellationToken + ) + .ConfigureAwait(false) + ) + { + // consume items until exception + } + }); + + var durations = collector.For("pulse.stream_query.duration"); + + using (Assert.Multiple()) + { + _ = await Assert.That(durations).Count().IsEqualTo(2); + _ = await Assert.That(durations[0].Unit).IsEqualTo("s"); + _ = await Assert.That(durations[0].Value).IsEqualTo(2d); + _ = await Assert.That(collector.For("pulse.stream_query.total")[0].Unit).IsEqualTo("{query}"); + _ = await Assert.That(collector.For("pulse.stream_query.errors")[0].Unit).IsEqualTo("{error}"); + } + } + + [Test] + [NotInParallel("PulseStreamQueryMetrics")] + public async Task HandleAsync_WithDefaultOptions_KeepsLegacyUnits(CancellationToken cancellationToken) + { + cancellationToken.ThrowIfCancellationRequested(); + + using var collector = new PulseMeasurementCollector("pulse.request.name", nameof(MeasuredStreamQuery)); + var interceptor = new ActivityAndMetricsStreamQueryInterceptor( + new FakeTimeProvider() + ); + + _ = await Assert.ThrowsAsync(async () => + { + await foreach ( + var _ in interceptor + .HandleAsync( + new MeasuredStreamQuery(), + (_, ct) => ThrowingItems(new InvalidOperationException("boom"), ct), + cancellationToken + ) + .ConfigureAwait(false) + ) + { + // consume items until exception + } + }); + + using (Assert.Multiple()) + { + _ = await Assert.That(collector.For("pulse.stream_query.duration")[0].Unit).IsEqualTo("ms"); + _ = await Assert.That(collector.For("pulse.stream_query.total")[0].Unit).IsEqualTo("queries"); + _ = await Assert.That(collector.For("pulse.stream_query.errors")[0].Unit).IsEqualTo("errors"); + } + } + private static IAsyncEnumerable Items(IEnumerable items, CancellationToken cancellationToken = default) => ItemsCore(items, cancellationToken); diff --git a/tests/NetEvolve.Pulse.Tests.Unit/Interceptors/PulseMeasurementCollector.cs b/tests/NetEvolve.Pulse.Tests.Unit/Interceptors/PulseMeasurementCollector.cs index 94793410..f0d8b628 100644 --- a/tests/NetEvolve.Pulse.Tests.Unit/Interceptors/PulseMeasurementCollector.cs +++ b/tests/NetEvolve.Pulse.Tests.Unit/Interceptors/PulseMeasurementCollector.cs @@ -8,16 +8,17 @@ namespace NetEvolve.Pulse.Tests.Unit.Interceptors; /// /// Collects measurements of the NetEvolve.Pulse meter that carry a specific tag value, so tests -/// running in parallel with other Pulse telemetry only see their own measurements. +/// running in parallel with other Pulse telemetry only see their own measurements. Without a tag filter, +/// all measurements are collected. /// internal sealed class PulseMeasurementCollector : IDisposable { private readonly MeterListener _listener = new(); private readonly ConcurrentQueue _measurements = new(); - private readonly string _tagKey; - private readonly string _tagValue; + private readonly string? _tagKey; + private readonly string? _tagValue; - public PulseMeasurementCollector(string tagKey, string tagValue) + public PulseMeasurementCollector(string? tagKey = null, string? tagValue = null) { _tagKey = tagKey; _tagValue = tagValue; @@ -48,8 +49,11 @@ private void Add(Instrument instrument, double value, ReadOnlySpan m.Unit == "s")).IsTrue(); + _ = await Assert + .That(collector.For("pulse.outbox.processed.total").Any(m => m.Unit == "{message}")) + .IsTrue(); + } + } + [Test] [NotInParallel("OutboxMetrics")] public async Task ExecuteAsync_WithPendingMessages_ObservableGaugeReflectsPendingCount( From 32b7987faea2f14a617bb3a5e33e28531ec32c80 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Martin=20St=C3=BChmer?= Date: Tue, 29 Sep 2026 01:11:38 +0200 Subject: [PATCH 06/14] test(outbox): wait for the processing duration measurement instead of counting polls --- .../Outbox/OutboxProcessorHostedServiceTests.cs | 9 ++++++++- 1 file changed, 8 insertions(+), 1 deletion(-) diff --git a/tests/NetEvolve.Pulse.Tests.Unit/Outbox/OutboxProcessorHostedServiceTests.cs b/tests/NetEvolve.Pulse.Tests.Unit/Outbox/OutboxProcessorHostedServiceTests.cs index 17e5b124..cc83ee8b 100644 --- a/tests/NetEvolve.Pulse.Tests.Unit/Outbox/OutboxProcessorHostedServiceTests.cs +++ b/tests/NetEvolve.Pulse.Tests.Unit/Outbox/OutboxProcessorHostedServiceTests.cs @@ -1180,7 +1180,14 @@ CancellationToken cancellationToken await service.StartAsync(cts.Token).ConfigureAwait(false); using var timeoutCts = CreateSignalTimeout(cancellationToken); await repository.WaitForMarkingsAsync(1, timeoutCts.Token).ConfigureAwait(false); - await repository.WaitForPollAsync(2, timeoutCts.Token).ConfigureAwait(false); + + // The duration is recorded after ProcessBatchAsync returns; the pending-count refresh also signals a + // poll, so wait for further polls until the measurement arrived instead of counting signals. + while (!collector.For("pulse.outbox.processing.duration").Any(m => m.Unit == "s")) + { + await repository.WaitForPollAsync(1, timeoutCts.Token).ConfigureAwait(false); + } + await cts.CancelAsync().ConfigureAwait(false); await service.StopAsync(cancellationToken).ConfigureAwait(false); From d3e5afc2294eba08b82ed1e4bcb0fe65036fdc9d Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Martin=20St=C3=BChmer?= Date: Tue, 29 Sep 2026 01:11:39 +0200 Subject: [PATCH 07/14] feat(telemetry): add opt-in OpenTelemetry semantic convention units for Pulse metrics Adds ActivityAndMetricsOptions.UseSemanticConventionUnits and an AddActivityAndMetrics(configure) overload. When enabled, the interceptor and outbox duration histograms record seconds (unit s) with explicit bucket boundaries, and the counters use UCUM annotation units ({request}, {event}, {query}, {error}, {message}). The default keeps the legacy units. The README documents the telemetry conventions and the migration steps. --- ...-28-opentelemetry-telemetry-conventions.md | 2 +- .../ActivityMetricsExtensions.cs | 30 ++++++ ...ivityAndMetricsEventInterceptor{TEvent}.cs | 75 ++++++++++----- .../Interceptors/ActivityAndMetricsOptions.cs | 27 ++++++ ...sRequestInterceptor{TRequest,TResponse}.cs | 75 ++++++++++----- ...treamQueryInterceptor{TQuery,TResponse}.cs | 86 +++++++++++------- .../Internals/TelemetryUnits.cs | 89 ++++++++++++++++++ .../Outbox/OutboxProcessorHostedService.cs | 91 +++++++++++-------- src/NetEvolve.Pulse/README.md | 35 +++++++ 9 files changed, 394 insertions(+), 116 deletions(-) create mode 100644 src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsOptions.cs create mode 100644 src/NetEvolve.Pulse/Internals/TelemetryUnits.cs diff --git a/decisions/2026-09-28-opentelemetry-telemetry-conventions.md b/decisions/2026-09-28-opentelemetry-telemetry-conventions.md index b6c06748..97a13ee4 100644 --- a/decisions/2026-09-28-opentelemetry-telemetry-conventions.md +++ b/decisions/2026-09-28-opentelemetry-telemetry-conventions.md @@ -6,7 +6,7 @@ applyTo: - "src/NetEvolve.Pulse/Interceptors/ActivityAndMetrics*.cs" - "src/NetEvolve.Pulse/Outbox/OutboxProcessorHostedService.cs" - "src/NetEvolve.Pulse/Internals/Defaults.cs" - - "src/NetEvolve.Pulse/Internals/OperationInstruments.cs" + - "src/NetEvolve.Pulse/Internals/TelemetryUnits.cs" created: 2026-09-28 diff --git a/src/NetEvolve.Pulse/ActivityMetricsExtensions.cs b/src/NetEvolve.Pulse/ActivityMetricsExtensions.cs index 223f9fe7..72e94e8e 100644 --- a/src/NetEvolve.Pulse/ActivityMetricsExtensions.cs +++ b/src/NetEvolve.Pulse/ActivityMetricsExtensions.cs @@ -19,11 +19,17 @@ public static class ActivityMetricsExtensions /// /// The mediator builder. /// The builder for method chaining. + /// + /// Metrics keep their legacy units by default. Use + /// to opt into the + /// units of the OpenTelemetry semantic conventions. + /// /// Thrown when is . public static IMediatorBuilder AddActivityAndMetrics(this IMediatorBuilder builder) { ArgumentNullException.ThrowIfNull(builder); + _ = builder.Services.AddOptions(); builder.Services.TryAddEnumerable( ServiceDescriptor.Singleton(typeof(IEventInterceptor<>), typeof(ActivityAndMetricsEventInterceptor<>)) ); @@ -39,4 +45,28 @@ public static IMediatorBuilder AddActivityAndMetrics(this IMediatorBuilder build return builder; } + + /// + /// Adds activity tracing and metrics collection for all requests processed by the mediator and configures + /// the telemetry options, for example to opt into the units of the OpenTelemetry semantic conventions. + /// The options also apply to the metrics of the outbox processor. + /// + /// The mediator builder. + /// The action that configures the . + /// The builder for method chaining. + /// + /// Thrown when or is . + /// + public static IMediatorBuilder AddActivityAndMetrics( + this IMediatorBuilder builder, + Action configure + ) + { + ArgumentNullException.ThrowIfNull(builder); + ArgumentNullException.ThrowIfNull(configure); + + _ = builder.Services.Configure(configure); + + return builder.AddActivityAndMetrics(); + } } diff --git a/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsEventInterceptor{TEvent}.cs b/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsEventInterceptor{TEvent}.cs index b36a4342..a0d202fb 100644 --- a/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsEventInterceptor{TEvent}.cs +++ b/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsEventInterceptor{TEvent}.cs @@ -2,6 +2,7 @@ using System.Diagnostics; using System.Diagnostics.Metrics; +using Microsoft.Extensions.Options; using NetEvolve.Pulse.Extensibility; using NetEvolve.Pulse.Internals; using static Internals.Defaults.Tags; @@ -16,31 +17,24 @@ internal sealed class ActivityAndMetricsEventInterceptor : IEventInterce where TEvent : IEvent { /// - /// Counter tracking the total number of events processed, tagged by event type. + /// Counter pulse.events.total, tagged by type. /// - private static readonly Counter EventCounter = Defaults.Meter.CreateCounter( - "pulse.events.total", - "events", - "Total number of events processed." - ); + private readonly Counter _eventCounter; /// - /// Counter tracking the total number of event errors, tagged by event type. + /// Counter pulse.event.errors, tagged by type. /// - private static readonly Counter ErrorsCounter = Defaults.Meter.CreateCounter( - "pulse.event.errors", - "errors", - "Total number of event errors." - ); + private readonly Counter _errorsCounter; /// - /// Histogram measuring event processing duration in milliseconds, with percentile distributions. + /// Histogram pulse.event.duration in milliseconds, or in seconds when semantic convention units are enabled. /// - private static readonly Histogram EventDurationHistogram = Defaults.Meter.CreateHistogram( - "pulse.event.duration", - "ms", - "Duration of event processing in milliseconds." - ); + private readonly Histogram _eventDurationHistogram; + + /// + /// Whether the metrics use the units of the OpenTelemetry semantic conventions. + /// + private readonly bool _useSemanticConventionUnits; /// /// Time provider for consistent timestamp generation, supporting testability. @@ -51,7 +45,37 @@ internal sealed class ActivityAndMetricsEventInterceptor : IEventInterce /// Initializes a new instance of the class. /// /// The time provider for timestamp generation. - public ActivityAndMetricsEventInterceptor(TimeProvider timeProvider) => _timeProvider = timeProvider; + /// The telemetry options; keeps the legacy units. + public ActivityAndMetricsEventInterceptor( + TimeProvider timeProvider, + IOptions? options = null + ) + { + _timeProvider = timeProvider; + _useSemanticConventionUnits = options?.Value.UseSemanticConventionUnits ?? false; + _eventCounter = TelemetryUnits.CreateCounter( + Defaults.Meter, + "pulse.events.total", + "events", + "{event}", + "Total number of events processed.", + _useSemanticConventionUnits + ); + _errorsCounter = TelemetryUnits.CreateCounter( + Defaults.Meter, + "pulse.event.errors", + "errors", + "{error}", + "Total number of event errors.", + _useSemanticConventionUnits + ); + _eventDurationHistogram = TelemetryUnits.CreateDurationHistogram( + Defaults.Meter, + "pulse.event.duration", + "event processing", + _useSemanticConventionUnits + ); + } /// /// @@ -95,7 +119,7 @@ public async Task HandleAsync( .SetTag(EventCorrelationId, message.CorrelationId) .SetTag(EventCausationId, message.CausationId) .SetTag(EventTimestamp, startTime); - EventCounter.Add(1, tags); + _eventCounter.Add(1, tags); try { @@ -111,7 +135,10 @@ public async Task HandleAsync( .SetTag(Success, value: true); // Record successful execution duration - EventDurationHistogram.Record((endTime - startTime).TotalMilliseconds, [.. tags, new(Success, true)]); + _eventDurationHistogram.Record( + TelemetryUnits.ToDuration(endTime - startTime, _useSemanticConventionUnits), + [.. tags, new(Success, true)] + ); } catch (Exception ex) { @@ -130,9 +157,9 @@ public async Task HandleAsync( .SetTag(Success, value: false); // Increment error counters and record failed execution duration - ErrorsCounter.Add(1, [.. tags, new(ErrorType, errorType)]); - EventDurationHistogram.Record( - (errorTime - startTime).TotalMilliseconds, + _errorsCounter.Add(1, [.. tags, new(ErrorType, errorType)]); + _eventDurationHistogram.Record( + TelemetryUnits.ToDuration(errorTime - startTime, _useSemanticConventionUnits), [.. tags, new(Success, false), new(ErrorType, errorType)] ); diff --git a/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsOptions.cs b/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsOptions.cs new file mode 100644 index 00000000..7c30e001 --- /dev/null +++ b/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsOptions.cs @@ -0,0 +1,27 @@ +namespace NetEvolve.Pulse.Interceptors; + +/// +/// Options for the built-in activity and metrics telemetry registered via AddActivityAndMetrics(). +/// The options also apply to the metrics of the outbox processor. +/// +/// +/// +/// services.AddPulse(c => c.AddActivityAndMetrics(o => o.UseSemanticConventionUnits = true)); +/// +/// +public sealed class ActivityAndMetricsOptions +{ + /// + /// Gets or sets a value indicating whether Pulse metrics use the units of the OpenTelemetry semantic conventions. + /// + /// + /// When , duration histograms record seconds with unit s and explicit bucket + /// boundaries, and counters use UCUM annotations ({request}, {event}, {query}, + /// {error}, {message}). When (default), durations are recorded in + /// milliseconds with unit ms and counters keep their legacy units (requests, events, + /// queries, errors, messages). Exporters such as Prometheus derive the exported metric + /// name from the unit, so enabling this option changes the exported names. A later 0.x release + /// makes the default. + /// + public bool UseSemanticConventionUnits { get; set; } +} diff --git a/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsRequestInterceptor{TRequest,TResponse}.cs b/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsRequestInterceptor{TRequest,TResponse}.cs index 081bd0a2..46112725 100644 --- a/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsRequestInterceptor{TRequest,TResponse}.cs +++ b/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsRequestInterceptor{TRequest,TResponse}.cs @@ -2,6 +2,7 @@ using System.Diagnostics; using System.Diagnostics.Metrics; +using Microsoft.Extensions.Options; using NetEvolve.Pulse.Extensibility; using NetEvolve.Pulse.Internals; using static Internals.Defaults.Tags; @@ -18,31 +19,24 @@ internal sealed class ActivityAndMetricsRequestInterceptor where TRequest : IRequest { /// - /// Counter tracking the total number of requests processed, tagged by request type. + /// Counter pulse.requests.total, tagged by type. /// - private static readonly Counter RequestCounter = Defaults.Meter.CreateCounter( - "pulse.requests.total", - "requests", - "Total number of requests processed." - ); + private readonly Counter _requestCounter; /// - /// Counter tracking the total number of request errors, tagged by request type. + /// Counter pulse.request.errors, tagged by type. /// - private static readonly Counter ErrorsCounter = Defaults.Meter.CreateCounter( - "pulse.request.errors", - "errors", - "Total number of request errors." - ); + private readonly Counter _errorsCounter; /// - /// Histogram measuring request processing duration in milliseconds, with percentile distributions. + /// Histogram pulse.request.duration in milliseconds, or in seconds when semantic convention units are enabled. /// - private static readonly Histogram RequestDurationHistogram = Defaults.Meter.CreateHistogram( - "pulse.request.duration", - "ms", - "Duration of request processing in milliseconds." - ); + private readonly Histogram _requestDurationHistogram; + + /// + /// Whether the metrics use the units of the OpenTelemetry semantic conventions. + /// + private readonly bool _useSemanticConventionUnits; /// /// Time provider for consistent timestamp generation, supporting testability. @@ -53,7 +47,37 @@ internal sealed class ActivityAndMetricsRequestInterceptor /// Initializes a new instance of the class. /// /// The time provider for timestamp generation. - public ActivityAndMetricsRequestInterceptor(TimeProvider timeProvider) => _timeProvider = timeProvider; + /// The telemetry options; keeps the legacy units. + public ActivityAndMetricsRequestInterceptor( + TimeProvider timeProvider, + IOptions? options = null + ) + { + _timeProvider = timeProvider; + _useSemanticConventionUnits = options?.Value.UseSemanticConventionUnits ?? false; + _requestCounter = TelemetryUnits.CreateCounter( + Defaults.Meter, + "pulse.requests.total", + "requests", + "{request}", + "Total number of requests processed.", + _useSemanticConventionUnits + ); + _errorsCounter = TelemetryUnits.CreateCounter( + Defaults.Meter, + "pulse.request.errors", + "errors", + "{error}", + "Total number of request errors.", + _useSemanticConventionUnits + ); + _requestDurationHistogram = TelemetryUnits.CreateDurationHistogram( + Defaults.Meter, + "pulse.request.duration", + "request processing", + _useSemanticConventionUnits + ); + } /// /// @@ -104,7 +128,7 @@ public async Task HandleAsync( .SetTag(RequestCorrelationId, request.CorrelationId) .SetTag(RequestCausationId, request.CausationId) .SetTag(RequestTimestamp, startTime); - RequestCounter.Add(1, tags); + _requestCounter.Add(1, tags); try { @@ -120,7 +144,10 @@ public async Task HandleAsync( .SetTag(Success, value: true); // Record successful execution duration - RequestDurationHistogram.Record((endTime - startTime).TotalMilliseconds, [.. tags, new(Success, true)]); + _requestDurationHistogram.Record( + TelemetryUnits.ToDuration(endTime - startTime, _useSemanticConventionUnits), + [.. tags, new(Success, true)] + ); return response; } @@ -141,9 +168,9 @@ public async Task HandleAsync( .SetTag(Success, value: false); // Increment error counters and record failed execution duration - ErrorsCounter.Add(1, [.. tags, new(ErrorType, errorType)]); - RequestDurationHistogram.Record( - (errorTime - startTime).TotalMilliseconds, + _errorsCounter.Add(1, [.. tags, new(ErrorType, errorType)]); + _requestDurationHistogram.Record( + TelemetryUnits.ToDuration(errorTime - startTime, _useSemanticConventionUnits), [.. tags, new(Success, false), new(ErrorType, errorType)] ); diff --git a/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsStreamQueryInterceptor{TQuery,TResponse}.cs b/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsStreamQueryInterceptor{TQuery,TResponse}.cs index c675065b..23ff5675 100644 --- a/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsStreamQueryInterceptor{TQuery,TResponse}.cs +++ b/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsStreamQueryInterceptor{TQuery,TResponse}.cs @@ -5,6 +5,7 @@ namespace NetEvolve.Pulse.Interceptors; using System.Diagnostics.Metrics; using System.Runtime.CompilerServices; using System.Runtime.ExceptionServices; +using Microsoft.Extensions.Options; using NetEvolve.Pulse.Extensibility; using NetEvolve.Pulse.Internals; using static Internals.Defaults.Tags; @@ -21,43 +22,36 @@ internal sealed class ActivityAndMetricsStreamQueryInterceptor { /// - /// Counter tracking the total number of stream queries processed, tagged by query type. + /// Cached query name derived from the generic type parameter. + /// Static fields in generic types are per type instantiation, so this is computed once per . /// - private static readonly Counter StreamQueryCounter = Defaults.Meter.CreateCounter( - "pulse.stream_query.total", - "queries", - "Total number of stream queries processed." - ); + private static readonly string QueryName = typeof(TQuery).Name; /// - /// Counter tracking the total number of stream query errors, tagged by query type. + /// Cached response type name derived from the generic type parameter. + /// Static fields in generic types are per type instantiation, so this is computed once per . /// - private static readonly Counter ErrorsCounter = Defaults.Meter.CreateCounter( - "pulse.stream_query.errors", - "errors", - "Total number of stream query errors." - ); + private static readonly string ResponseTypeName = typeof(TResponse).Name; /// - /// Histogram measuring stream query processing duration in milliseconds, with percentile distributions. + /// Counter pulse.stream_query.total, tagged by type. /// - private static readonly Histogram StreamQueryDurationHistogram = Defaults.Meter.CreateHistogram( - "pulse.stream_query.duration", - "ms", - "Duration of stream query processing in milliseconds." - ); + private readonly Counter _streamQueryCounter; /// - /// Cached query name derived from the generic type parameter. - /// Static fields in generic types are per type instantiation, so this is computed once per . + /// Counter pulse.stream_query.errors, tagged by type. /// - private static readonly string QueryName = typeof(TQuery).Name; + private readonly Counter _errorsCounter; /// - /// Cached response type name derived from the generic type parameter. - /// Static fields in generic types are per type instantiation, so this is computed once per . + /// Histogram pulse.stream_query.duration in milliseconds, or in seconds when semantic convention units are enabled. /// - private static readonly string ResponseTypeName = typeof(TResponse).Name; + private readonly Histogram _streamQueryDurationHistogram; + + /// + /// Whether the metrics use the units of the OpenTelemetry semantic conventions. + /// + private readonly bool _useSemanticConventionUnits; /// /// Time provider for consistent timestamp generation, supporting testability. @@ -68,7 +62,37 @@ internal sealed class ActivityAndMetricsStreamQueryInterceptor class. /// /// The time provider for timestamp generation. - public ActivityAndMetricsStreamQueryInterceptor(TimeProvider timeProvider) => _timeProvider = timeProvider; + /// The telemetry options; keeps the legacy units. + public ActivityAndMetricsStreamQueryInterceptor( + TimeProvider timeProvider, + IOptions? options = null + ) + { + _timeProvider = timeProvider; + _useSemanticConventionUnits = options?.Value.UseSemanticConventionUnits ?? false; + _streamQueryCounter = TelemetryUnits.CreateCounter( + Defaults.Meter, + "pulse.stream_query.total", + "queries", + "{query}", + "Total number of stream queries processed.", + _useSemanticConventionUnits + ); + _errorsCounter = TelemetryUnits.CreateCounter( + Defaults.Meter, + "pulse.stream_query.errors", + "errors", + "{error}", + "Total number of stream query errors.", + _useSemanticConventionUnits + ); + _streamQueryDurationHistogram = TelemetryUnits.CreateDurationHistogram( + Defaults.Meter, + "pulse.stream_query.duration", + "stream query processing", + _useSemanticConventionUnits + ); + } /// /// @@ -122,7 +146,7 @@ public async IAsyncEnumerable HandleAsync( .SetTag(RequestCorrelationId, request.CorrelationId) .SetTag(StreamQueryCausationId, request.CausationId) .SetTag(RequestTimestamp, startTime); - StreamQueryCounter.Add(1, tags); + _streamQueryCounter.Add(1, tags); // yield return is not allowed inside a try/catch block, so we capture any exception // from the handler or the inner enumerator and re-throw it after the finally block ran. @@ -200,7 +224,7 @@ bool completed ) { var endTime = _timeProvider.GetUtcNow(); - var duration = (endTime - startTime).TotalMilliseconds; + var duration = TelemetryUnits.ToDuration(endTime - startTime, _useSemanticConventionUnits); _ = activity?.SetEndTime(endTime.UtcDateTime); if (exception is not null) @@ -217,22 +241,22 @@ bool completed .SetTag(ExceptionTimestamp, endTime) .SetTag(Success, value: false); - ErrorsCounter.Add(1, [.. tags, new(ErrorType, errorType)]); - StreamQueryDurationHistogram.Record(duration, [.. tags, new(Success, false), new(ErrorType, errorType)]); + _errorsCounter.Add(1, [.. tags, new(ErrorType, errorType)]); + _streamQueryDurationHistogram.Record(duration, [.. tags, new(Success, false), new(ErrorType, errorType)]); } else if (completed) { // Status stays Unset on success, as required by the OpenTelemetry Trace API. _ = activity?.SetTag(ResponseTimestamp, endTime).SetTag(Success, value: true); - StreamQueryDurationHistogram.Record(duration, [.. tags, new(Success, true)]); + _streamQueryDurationHistogram.Record(duration, [.. tags, new(Success, true)]); } else { // The consumer stopped early: not an error, so the status stays Unset. _ = activity?.SetTag(StreamCompleted, value: false); - StreamQueryDurationHistogram.Record(duration, [.. tags, new(StreamCompleted, false)]); + _streamQueryDurationHistogram.Record(duration, [.. tags, new(StreamCompleted, false)]); } } } diff --git a/src/NetEvolve.Pulse/Internals/TelemetryUnits.cs b/src/NetEvolve.Pulse/Internals/TelemetryUnits.cs new file mode 100644 index 00000000..0c47f341 --- /dev/null +++ b/src/NetEvolve.Pulse/Internals/TelemetryUnits.cs @@ -0,0 +1,89 @@ +namespace NetEvolve.Pulse.Internals; + +using System.Diagnostics.Metrics; + +/// +/// Creates Pulse instruments with either the legacy units or the units of the OpenTelemetry semantic conventions, +/// as selected by . +/// +internal static class TelemetryUnits +{ + /// + /// Bucket boundaries in seconds for duration histograms, taken from the OpenTelemetry HTTP semantic conventions. + /// + private static readonly double[] DurationBucketBoundariesInSeconds = + [ + 0.005, + 0.01, + 0.025, + 0.05, + 0.075, + 0.1, + 0.25, + 0.5, + 0.75, + 1, + 2.5, + 5, + 7.5, + 10, + ]; + + /// + /// Creates a counter on with the legacy or the UCUM annotation unit. + /// + /// The meter that owns the counter. + /// The instrument name. + /// The legacy unit, for example requests. + /// The UCUM annotation unit, for example {request}. + /// The instrument description. + /// Whether to use the semantic convention unit. + /// The created counter. + public static Counter CreateCounter( + Meter meter, + string name, + string legacyUnit, + string annotationUnit, + string description, + bool useSemanticConventionUnits + ) => meter.CreateCounter(name, useSemanticConventionUnits ? annotationUnit : legacyUnit, description); + + /// + /// Creates a duration histogram on in milliseconds (ms) or seconds (s). + /// + /// The meter that owns the histogram. + /// The instrument name. + /// The measured operation used in the description, for example request processing. + /// Whether to record seconds. + /// The created histogram. + public static Histogram CreateDurationHistogram( + Meter meter, + string name, + string subject, + bool useSemanticConventionUnits + ) + { + if (!useSemanticConventionUnits) + { + return meter.CreateHistogram(name, "ms", $"Duration of {subject} in milliseconds."); + } + + return meter.CreateHistogram( + name, + "s", + $"Duration of {subject} in seconds.", + tags: null, + advice: new InstrumentAdvice { HistogramBucketBoundaries = DurationBucketBoundariesInSeconds } + ); + } + + /// + /// Converts into the value recorded by a histogram created by + /// . + /// + /// The measured duration. + /// Whether the histogram records seconds. + /// The duration in seconds or milliseconds. + public static double ToDuration(TimeSpan elapsed, bool useSemanticConventionUnits) => + useSemanticConventionUnits ? elapsed.TotalSeconds : elapsed.TotalMilliseconds; +} diff --git a/src/NetEvolve.Pulse/Outbox/OutboxProcessorHostedService.cs b/src/NetEvolve.Pulse/Outbox/OutboxProcessorHostedService.cs index 94b74bc4..2e2f76d1 100644 --- a/src/NetEvolve.Pulse/Outbox/OutboxProcessorHostedService.cs +++ b/src/NetEvolve.Pulse/Outbox/OutboxProcessorHostedService.cs @@ -7,6 +7,7 @@ using Microsoft.Extensions.Logging; using Microsoft.Extensions.Options; using NetEvolve.Pulse.Extensibility.Outbox; +using NetEvolve.Pulse.Interceptors; using NetEvolve.Pulse.Internals; /// @@ -36,35 +37,23 @@ /// internal sealed partial class OutboxProcessorHostedService : BackgroundService { -#pragma warning disable IDE1006 // Naming rule violation - matching existing static metric field naming conventions - /// Counter tracking the total number of successfully processed outbox messages. - private static readonly Counter ProcessedCounter = Defaults.Meter.CreateCounter( - "pulse.outbox.processed.total", - "messages", - "Cumulative number of successfully processed outbox messages." - ); + /// Counter pulse.outbox.processed.total of successfully processed outbox messages. + private readonly Counter _processedCounter; - /// Counter tracking the total number of failed outbox processing attempts. - private static readonly Counter FailedCounter = Defaults.Meter.CreateCounter( - "pulse.outbox.failed.total", - "messages", - "Cumulative number of failed outbox processing attempts." - ); + /// Counter pulse.outbox.failed.total of failed outbox processing attempts. + private readonly Counter _failedCounter; - /// Counter tracking the total number of outbox messages moved to dead-letter. - private static readonly Counter DeadLetterCounter = Defaults.Meter.CreateCounter( - "pulse.outbox.deadletter.total", - "messages", - "Cumulative number of outbox messages moved to dead-letter." - ); + /// Counter pulse.outbox.deadletter.total of outbox messages moved to dead-letter. + private readonly Counter _deadLetterCounter; - /// Histogram measuring the duration of each outbox processing batch in milliseconds. - private static readonly Histogram ProcessingDurationHistogram = Defaults.Meter.CreateHistogram( - "pulse.outbox.processing.duration", - "ms", - "Duration of each outbox processing batch in milliseconds." - ); -#pragma warning restore IDE1006 + /// + /// Histogram pulse.outbox.processing.duration of each processing batch, in milliseconds or, with + /// semantic convention units, in seconds. + /// + private readonly Histogram _processingDurationHistogram; + + /// Whether the metrics use the units of the OpenTelemetry semantic conventions. + private readonly bool _useSemanticConventionUnits; /// /// Creates the per-cycle and per-work-item scopes used to resolve the scoped @@ -113,13 +102,15 @@ internal sealed partial class OutboxProcessorHostedService : BackgroundService /// The processor configuration options. /// The logger for diagnostic output. /// The time provider used for retry scheduling and the polling delays. + /// The telemetry options; keeps the legacy metric units. public OutboxProcessorHostedService( IServiceScopeFactory scopeFactory, IMessageTransport transport, IHostApplicationLifetime lifetime, IOptions options, ILogger logger, - TimeProvider timeProvider + TimeProvider timeProvider, + IOptions? telemetryOptions = null ) { ArgumentNullException.ThrowIfNull(scopeFactory); @@ -136,11 +127,36 @@ TimeProvider timeProvider _logger = logger; _timeProvider = timeProvider; + _useSemanticConventionUnits = telemetryOptions?.Value.UseSemanticConventionUnits ?? false; + var messageUnit = _useSemanticConventionUnits ? "{message}" : "messages"; + + _processedCounter = Defaults.Meter.CreateCounter( + "pulse.outbox.processed.total", + messageUnit, + "Cumulative number of successfully processed outbox messages." + ); + _failedCounter = Defaults.Meter.CreateCounter( + "pulse.outbox.failed.total", + messageUnit, + "Cumulative number of failed outbox processing attempts." + ); + _deadLetterCounter = Defaults.Meter.CreateCounter( + "pulse.outbox.deadletter.total", + messageUnit, + "Cumulative number of outbox messages moved to dead-letter." + ); + _processingDurationHistogram = TelemetryUnits.CreateDurationHistogram( + Defaults.Meter, + "pulse.outbox.processing.duration", + "each outbox processing batch", + _useSemanticConventionUnits + ); + _meter = new Meter(Defaults.Meter.Name, Defaults.Version); _ = _meter.CreateObservableGauge( "pulse.outbox.pending", observeValue: () => Volatile.Read(ref _pendingCount), - unit: "messages", + unit: messageUnit, description: "Current number of pending outbox messages." ); } @@ -207,11 +223,14 @@ protected override async Task ExecuteAsync(CancellationToken stoppingToken) var batchStartTime = Stopwatch.GetTimestamp(); var succeededCount = await ProcessBatchAsync(repository, stoppingToken).ConfigureAwait(false); - var elapsed = Stopwatch.GetElapsedTime(batchStartTime).TotalMilliseconds; + var elapsed = TelemetryUnits.ToDuration( + Stopwatch.GetElapsedTime(batchStartTime), + _useSemanticConventionUnits + ); try { - ProcessingDurationHistogram.Record(elapsed); + _processingDurationHistogram.Record(elapsed); } catch (Exception ex) { @@ -440,7 +459,7 @@ CancellationToken cancellationToken try { - ProcessedCounter.Add(1); + _processedCounter.Add(1); } catch (Exception ex) { @@ -464,7 +483,7 @@ CancellationToken cancellationToken try { - DeadLetterCounter.Add(1); + _deadLetterCounter.Add(1); } catch (Exception metricEx) { @@ -481,7 +500,7 @@ await repository try { - FailedCounter.Add(1); + _failedCounter.Add(1); } catch (Exception metricEx) { @@ -555,7 +574,7 @@ CancellationToken cancellationToken try { - ProcessedCounter.Add(messages.Length); + _processedCounter.Add(messages.Length); } catch (Exception ex) { @@ -626,12 +645,12 @@ await Task.WhenAll( { if (failedMessages.Length > 0) { - FailedCounter.Add(failedMessages.Length); + _failedCounter.Add(failedMessages.Length); } if (deadLetterMessages.Length > 0) { - DeadLetterCounter.Add(deadLetterMessages.Length); + _deadLetterCounter.Add(deadLetterMessages.Length); } } catch (Exception metricEx) diff --git a/src/NetEvolve.Pulse/README.md b/src/NetEvolve.Pulse/README.md index 53eaf0ae..e1f80357 100644 --- a/src/NetEvolve.Pulse/README.md +++ b/src/NetEvolve.Pulse/README.md @@ -325,6 +325,41 @@ public sealed class NewtonsoftJsonPayloadSerializer : IPayloadSerializer The custom serializer will be used for all payload operations within Pulse. Ensure your implementation is thread-safe, as the same instance may be accessed concurrently from multiple pipeline stages. +## Telemetry + +`AddActivityAndMetrics()` emits activities and metrics on the `NetEvolve.Pulse` activity source and meter, following the [OpenTelemetry recording errors conventions](https://opentelemetry.io/docs/specs/semconv/general/recording-errors/): + +* Successful operations leave the activity status `Unset`. Failed operations set `Error` with the exception message. +* Failed operations carry `error.type` (the full exception type name) on the activity, the error counter and the duration histogram. +* A stream query records its duration once for every outcome. A stream whose consumer stops early (`break`, `Take`, a disconnected client) carries `pulse.stream.completed=false` instead of `pulse.success`. + +### Semantic Convention Units + +Metrics keep their legacy units by default. Opt into the units of the [OpenTelemetry metrics guidelines](https://opentelemetry.io/docs/specs/semconv/general/metrics/#instrument-units) once your dashboards and alerts are migrated: + +```csharp +services.AddPulse(config => config.AddActivityAndMetrics(options => options.UseSemanticConventionUnits = true)); +``` + +The option also applies to the outbox processor metrics. Without `AddActivityAndMetrics`, configure it with `services.Configure(...)`. + +| Instrument | Default unit | With `UseSemanticConventionUnits` | +| --- | --- | --- | +| `pulse.requests.total` | `requests` | `{request}` | +| `pulse.events.total` | `events` | `{event}` | +| `pulse.stream_query.total` | `queries` | `{query}` | +| `pulse.request.errors`, `pulse.event.errors`, `pulse.stream_query.errors` | `errors` | `{error}` | +| `pulse.outbox.processed.total`, `pulse.outbox.failed.total`, `pulse.outbox.deadletter.total`, `pulse.outbox.pending` | `messages` | `{message}` | +| `pulse.request.duration`, `pulse.event.duration`, `pulse.stream_query.duration`, `pulse.outbox.processing.duration` | `ms` (milliseconds) | `s` (seconds, with bucket boundaries from 0.005 to 10 s) | + +Migration steps: + +1. Replace filters on the `Ok` span status with "status is not `Error`". +2. Enable `UseSemanticConventionUnits`. Exporters that append the unit to the metric name (for example the Prometheus exporter) then export new names such as `pulse_request_duration_seconds` instead of `pulse_request_duration_milliseconds`. +3. Update dashboards and alerts to the new names and convert duration thresholds from milliseconds to seconds. + +A later `0.x` release makes the semantic convention units the default. + ## NativeAOT and Trimming All Pulse runtime packages are built with `IsAotCompatible` enabled, so the trim and NativeAOT analyzers run on every build. The core mediator pipeline (`AddPulse`, handler registration through `NetEvolve.Pulse.SourceGeneration` or the generic `Add*Handler<,>` methods, `SendAsync`, `QueryAsync`, `StreamQueryAsync` and `PublishAsync`) is trim- and NativeAOT-safe. This is verified on every pull request by publishing and running the `samples/NetEvolve.Pulse.Xample.Aot` smoke application with NativeAOT for `net8.0`, `net9.0` and `net10.0`. It covers commands, queries, stream queries and events, open-generic handlers and interceptors, value-type and `Void` requests through the built-in interceptors, the outbox event type round-trip and the payload serializer setup below. From 5c0381c36e5cd912775197c127e28408112d2611 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Martin=20St=C3=BChmer?= Date: Tue, 29 Sep 2026 01:35:49 +0200 Subject: [PATCH 08/14] docs(decisions): apply duration bucket advice on all target frameworks --- decisions/2026-09-28-opentelemetry-telemetry-conventions.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/decisions/2026-09-28-opentelemetry-telemetry-conventions.md b/decisions/2026-09-28-opentelemetry-telemetry-conventions.md index 97a13ee4..11d84716 100644 --- a/decisions/2026-09-28-opentelemetry-telemetry-conventions.md +++ b/decisions/2026-09-28-opentelemetry-telemetry-conventions.md @@ -55,7 +55,7 @@ No earlier decision defines the Pulse telemetry contract. Units are part of the | `pulse.outbox.processed.total`, `pulse.outbox.failed.total`, `pulse.outbox.deadletter.total`, `pulse.outbox.pending` | `messages` | `{message}` | | `pulse.request.duration`, `pulse.event.duration`, `pulse.stream_query.duration`, `pulse.outbox.processing.duration` | `ms` | `s` | -* MUST provide explicit histogram bucket boundaries for the `s` unit on .NET 9 and later (`InstrumentAdvice`), using the boundaries of the HTTP semantic conventions: 0.005, 0.01, 0.025, 0.05, 0.075, 0.1, 0.25, 0.5, 0.75, 1, 2.5, 5, 7.5, 10. +* MUST provide explicit histogram bucket boundaries for the `s` unit on all target frameworks (`InstrumentAdvice`), using the boundaries of the HTTP semantic conventions: 0.005, 0.01, 0.025, 0.05, 0.075, 0.1, 0.25, 0.5, 0.75, 1, 2.5, 5, 7.5, 10. * MUST keep the metric names and tag names unchanged. * Transition plan: a later `0.x` release makes `UseSemanticConventionUnits` the default. The legacy units and the option are removed before `1.0.0`. From 0e96031235676e47a83988f7a8c9a69c9542f57a Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Martin=20St=C3=BChmer?= Date: Tue, 29 Sep 2026 02:11:02 +0200 Subject: [PATCH 09/14] test(interceptors): cover throwing inner dispose, throwing Current and cancellation in the stream telemetry interceptor --- ...tyAndMetricsStreamQueryInterceptorTests.cs | 241 ++++++++++++++++++ 1 file changed, 241 insertions(+) diff --git a/tests/NetEvolve.Pulse.Tests.Unit/Interceptors/ActivityAndMetricsStreamQueryInterceptorTests.cs b/tests/NetEvolve.Pulse.Tests.Unit/Interceptors/ActivityAndMetricsStreamQueryInterceptorTests.cs index 7475715a..142103f3 100644 --- a/tests/NetEvolve.Pulse.Tests.Unit/Interceptors/ActivityAndMetricsStreamQueryInterceptorTests.cs +++ b/tests/NetEvolve.Pulse.Tests.Unit/Interceptors/ActivityAndMetricsStreamQueryInterceptorTests.cs @@ -1,6 +1,8 @@ namespace NetEvolve.Pulse.Tests.Unit.Interceptors; +using System.Collections.Concurrent; using System.Diagnostics; +using System.Diagnostics.Metrics; using System.Runtime.CompilerServices; using Microsoft.Extensions.Options; using Microsoft.Extensions.Time.Testing; @@ -671,6 +673,214 @@ var _ in interceptor } } + [Test] + [NotInParallel("PulseStreamQueryMetrics")] + public async Task HandleAsync_WhenDisposeThrowsAfterCompletion_RecordsDisposeFailure( + CancellationToken cancellationToken + ) + { + cancellationToken.ThrowIfCancellationRequested(); + + using var collector = new PulseMeasurementCollector("pulse.request.name", nameof(MeasuredStreamQuery)); + var interceptor = new ActivityAndMetricsStreamQueryInterceptor(TimeProvider.System); + var stream = new FaultyStream(hasItem: false, disposeException: new ObjectDisposedException("inner")); + + _ = await Assert.ThrowsAsync(async () => + { + await foreach ( + var _ in interceptor + .HandleAsync(new MeasuredStreamQuery(), (_, _) => stream, cancellationToken) + .ConfigureAwait(false) + ) + { + // no items expected + } + }); + + var durations = collector.For("pulse.stream_query.duration"); + var errors = collector.For("pulse.stream_query.errors"); + + using (Assert.Multiple()) + { + _ = await Assert.That(durations).Count().IsEqualTo(1); + _ = await Assert.That(durations[0].Tags["pulse.success"] is false).IsTrue(); + _ = await Assert.That(durations[0].Tags["error.type"]).IsEqualTo("System.ObjectDisposedException"); + _ = await Assert.That(errors).Count().IsEqualTo(1); + } + } + + [Test] + [NotInParallel("PulseStreamQueryMetrics")] + public async Task HandleAsync_WhenDisposeThrowsAfterFault_ThrowsTheRecordedException( + CancellationToken cancellationToken + ) + { + cancellationToken.ThrowIfCancellationRequested(); + + using var collector = new PulseMeasurementCollector("pulse.request.name", nameof(MeasuredStreamQuery)); + var interceptor = new ActivityAndMetricsStreamQueryInterceptor(TimeProvider.System); + var stream = new FaultyStream( + moveNextException: new InvalidOperationException("fault"), + disposeException: new ObjectDisposedException("inner") + ); + + var exception = await Assert.ThrowsAsync(async () => + { + await foreach ( + var _ in interceptor + .HandleAsync(new MeasuredStreamQuery(), (_, _) => stream, cancellationToken) + .ConfigureAwait(false) + ) + { + // no items expected + } + }); + + var durations = collector.For("pulse.stream_query.duration"); + + using (Assert.Multiple()) + { + _ = await Assert.That(durations).Count().IsEqualTo(1); + _ = await Assert.That(durations[0].Tags["error.type"]).IsEqualTo(exception!.GetType().FullName); + } + } + + [Test] + [NotInParallel("PulseStreamQueryMetrics")] + public async Task HandleAsync_WhenDisposeThrowsAfterEarlyStop_RecordsDisposeFailure( + CancellationToken cancellationToken + ) + { + cancellationToken.ThrowIfCancellationRequested(); + + using var collector = new PulseMeasurementCollector("pulse.request.name", nameof(MeasuredStreamQuery)); + var interceptor = new ActivityAndMetricsStreamQueryInterceptor(TimeProvider.System); + var stream = new FaultyStream(disposeException: new ObjectDisposedException("inner")); + var consumed = 0; + + _ = await Assert.ThrowsAsync(async () => + { + await foreach ( + var _ in interceptor + .HandleAsync(new MeasuredStreamQuery(), (_, _) => stream, cancellationToken) + .ConfigureAwait(false) + ) + { + consumed++; + if (consumed == 1) + { + break; + } + } + }); + + var durations = collector.For("pulse.stream_query.duration"); + + using (Assert.Multiple()) + { + _ = await Assert.That(durations).Count().IsEqualTo(1); + _ = await Assert.That(durations[0].Tags["error.type"]).IsEqualTo("System.ObjectDisposedException"); + _ = await Assert.That(durations[0].Tags.ContainsKey("pulse.stream.completed")).IsFalse(); + _ = await Assert.That(collector.For("pulse.stream_query.errors")).Count().IsEqualTo(1); + } + } + + [Test] + [NotInParallel("PulseStreamQueryMetrics")] + public async Task HandleAsync_WhenCurrentThrows_RecordsError(CancellationToken cancellationToken) + { + cancellationToken.ThrowIfCancellationRequested(); + + using var collector = new PulseMeasurementCollector("pulse.request.name", nameof(MeasuredStreamQuery)); + var interceptor = new ActivityAndMetricsStreamQueryInterceptor(TimeProvider.System); + var stream = new FaultyStream(currentException: new InvalidOperationException("current")); + + _ = await Assert.ThrowsAsync(async () => + { + await foreach ( + var _ in interceptor + .HandleAsync(new MeasuredStreamQuery(), (_, _) => stream, cancellationToken) + .ConfigureAwait(false) + ) + { + // no items expected + } + }); + + var durations = collector.For("pulse.stream_query.duration"); + + using (Assert.Multiple()) + { + _ = await Assert.That(durations).Count().IsEqualTo(1); + _ = await Assert.That(durations[0].Tags["error.type"]).IsEqualTo("System.InvalidOperationException"); + _ = await Assert.That(collector.For("pulse.stream_query.errors")).Count().IsEqualTo(1); + } + } + + [Test] + [NotInParallel("PulseStreamQueryMetrics")] + public async Task HandleAsync_WhenTokenIsCancelledDuringEnumeration_RecordsCancellationAsError( + CancellationToken cancellationToken + ) + { + cancellationToken.ThrowIfCancellationRequested(); + + using var collector = new PulseMeasurementCollector("pulse.request.name", nameof(MeasuredStreamQuery)); + var interceptor = new ActivityAndMetricsStreamQueryInterceptor(TimeProvider.System); + using var cts = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken); + + _ = await Assert.ThrowsAsync(async () => + { + await foreach ( + var _ in interceptor + .HandleAsync(new MeasuredStreamQuery(), (_, ct) => Items([1, 2, 3], ct), cts.Token) + .ConfigureAwait(false) + ) + { + await cts.CancelAsync().ConfigureAwait(false); + } + }); + + var durations = collector.For("pulse.stream_query.duration"); + + using (Assert.Multiple()) + { + _ = await Assert.That(durations).Count().IsEqualTo(1); + _ = await Assert.That(durations[0].Tags["pulse.success"] is false).IsTrue(); + _ = await Assert.That(durations[0].Tags["error.type"]).IsEqualTo("System.OperationCanceledException"); + _ = await Assert.That(collector.For("pulse.stream_query.errors")).Count().IsEqualTo(1); + } + } + + [Test] + public async Task Constructor_CalledRepeatedly_ReusesInstruments(CancellationToken cancellationToken) + { + cancellationToken.ThrowIfCancellationRequested(); + + var published = new ConcurrentBag(); + using var listener = new MeterListener + { + InstrumentPublished = (instrument, _) => + { + if ( + string.Equals(instrument.Meter.Name, "NetEvolve.Pulse", StringComparison.Ordinal) + && string.Equals(instrument.Name, "pulse.stream_query.total", StringComparison.Ordinal) + && string.Equals(instrument.Unit, "queries", StringComparison.Ordinal) + ) + { + published.Add(instrument); + } + }, + }; + listener.Start(); + + _ = new ActivityAndMetricsStreamQueryInterceptor(TimeProvider.System); + _ = new ActivityAndMetricsStreamQueryInterceptor(TimeProvider.System); + _ = new ActivityAndMetricsStreamQueryInterceptor(TimeProvider.System); + + _ = await Assert.That(published.Distinct().Count()).IsEqualTo(1); + } + private static IAsyncEnumerable Items(IEnumerable items, CancellationToken cancellationToken = default) => ItemsCore(items, cancellationToken); @@ -698,6 +908,37 @@ private static async IAsyncEnumerable ThrowingItems( throw exception; } + /// + /// Stream whose enumerator fails at a chosen point. TUnit.Mocks cannot generate + /// on .NET 9 and later (CS9244 on the allows ref struct type parameter), so this is hand-written. + /// + private sealed class FaultyStream( + bool hasItem = true, + Exception? moveNextException = null, + Exception? currentException = null, + Exception? disposeException = null + ) : IAsyncEnumerable + { + public IAsyncEnumerator GetAsyncEnumerator(CancellationToken cancellationToken = default) => + new FaultyEnumerator(hasItem, moveNextException, currentException, disposeException); + + private sealed class FaultyEnumerator( + bool hasItem, + Exception? moveNextException, + Exception? currentException, + Exception? disposeException + ) : IAsyncEnumerator + { + public int Current => currentException is null ? 1 : throw currentException; + + public ValueTask MoveNextAsync() => + moveNextException is null ? ValueTask.FromResult(hasItem) : throw moveNextException; + + public ValueTask DisposeAsync() => + disposeException is null ? ValueTask.CompletedTask : throw disposeException; + } + } + private sealed class TestStreamQuery : IStreamQuery { public string? CausationId { get; set; } From 43f2338e1120c79a621eb6a7a2378abcbdcc362c Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Martin=20St=C3=BChmer?= Date: Tue, 29 Sep 2026 02:11:02 +0200 Subject: [PATCH 10/14] test(telemetry): require shared interceptor instruments and released outbox instruments --- ...ActivityAndMetricsEventInterceptorTests.cs | 31 ++++++++++++ ...tivityAndMetricsRequestInterceptorTests.cs | 31 ++++++++++++ .../OutboxProcessorHostedServiceTests.cs | 50 +++++++++++++++++++ 3 files changed, 112 insertions(+) diff --git a/tests/NetEvolve.Pulse.Tests.Unit/Interceptors/ActivityAndMetricsEventInterceptorTests.cs b/tests/NetEvolve.Pulse.Tests.Unit/Interceptors/ActivityAndMetricsEventInterceptorTests.cs index 10cc8462..4b62d3c3 100644 --- a/tests/NetEvolve.Pulse.Tests.Unit/Interceptors/ActivityAndMetricsEventInterceptorTests.cs +++ b/tests/NetEvolve.Pulse.Tests.Unit/Interceptors/ActivityAndMetricsEventInterceptorTests.cs @@ -1,6 +1,8 @@ namespace NetEvolve.Pulse.Tests.Unit.Interceptors; +using System.Collections.Concurrent; using System.Diagnostics; +using System.Diagnostics.Metrics; using Microsoft.Extensions.Options; using Microsoft.Extensions.Time.Testing; using NetEvolve.Extensions.TUnit; @@ -450,6 +452,35 @@ await interceptor } } + [Test] + public async Task Constructor_CalledRepeatedly_ReusesInstruments(CancellationToken cancellationToken) + { + cancellationToken.ThrowIfCancellationRequested(); + + var published = new ConcurrentBag(); + using var listener = new MeterListener + { + InstrumentPublished = (instrument, _) => + { + if ( + string.Equals(instrument.Meter.Name, "NetEvolve.Pulse", StringComparison.Ordinal) + && string.Equals(instrument.Name, "pulse.events.total", StringComparison.Ordinal) + && string.Equals(instrument.Unit, "events", StringComparison.Ordinal) + ) + { + published.Add(instrument); + } + }, + }; + listener.Start(); + + _ = new ActivityAndMetricsEventInterceptor(TimeProvider.System); + _ = new ActivityAndMetricsEventInterceptor(TimeProvider.System); + _ = new ActivityAndMetricsEventInterceptor(TimeProvider.System); + + _ = await Assert.That(published.Distinct().Count()).IsEqualTo(1); + } + private sealed class MeasuredEvent : IEvent { public string Id { get; init; } = Guid.NewGuid().ToString(); diff --git a/tests/NetEvolve.Pulse.Tests.Unit/Interceptors/ActivityAndMetricsRequestInterceptorTests.cs b/tests/NetEvolve.Pulse.Tests.Unit/Interceptors/ActivityAndMetricsRequestInterceptorTests.cs index 0f5ef228..96a58ceb 100644 --- a/tests/NetEvolve.Pulse.Tests.Unit/Interceptors/ActivityAndMetricsRequestInterceptorTests.cs +++ b/tests/NetEvolve.Pulse.Tests.Unit/Interceptors/ActivityAndMetricsRequestInterceptorTests.cs @@ -1,6 +1,8 @@ namespace NetEvolve.Pulse.Tests.Unit.Interceptors; +using System.Collections.Concurrent; using System.Diagnostics; +using System.Diagnostics.Metrics; using Microsoft.Extensions.Options; using Microsoft.Extensions.Time.Testing; using NetEvolve.Extensions.TUnit; @@ -508,6 +510,35 @@ await interceptor } } + [Test] + public async Task Constructor_CalledRepeatedly_ReusesInstruments(CancellationToken cancellationToken) + { + cancellationToken.ThrowIfCancellationRequested(); + + var published = new ConcurrentBag(); + using var listener = new MeterListener + { + InstrumentPublished = (instrument, _) => + { + if ( + string.Equals(instrument.Meter.Name, "NetEvolve.Pulse", StringComparison.Ordinal) + && string.Equals(instrument.Name, "pulse.requests.total", StringComparison.Ordinal) + && string.Equals(instrument.Unit, "requests", StringComparison.Ordinal) + ) + { + published.Add(instrument); + } + }, + }; + listener.Start(); + + _ = new ActivityAndMetricsRequestInterceptor(TimeProvider.System); + _ = new ActivityAndMetricsRequestInterceptor(TimeProvider.System); + _ = new ActivityAndMetricsRequestInterceptor(TimeProvider.System); + + _ = await Assert.That(published.Distinct().Count()).IsEqualTo(1); + } + private sealed class MeasuredCommand : ICommand { public string? CausationId { get; set; } diff --git a/tests/NetEvolve.Pulse.Tests.Unit/Outbox/OutboxProcessorHostedServiceTests.cs b/tests/NetEvolve.Pulse.Tests.Unit/Outbox/OutboxProcessorHostedServiceTests.cs index 0d509fe8..9aebd881 100644 --- a/tests/NetEvolve.Pulse.Tests.Unit/Outbox/OutboxProcessorHostedServiceTests.cs +++ b/tests/NetEvolve.Pulse.Tests.Unit/Outbox/OutboxProcessorHostedServiceTests.cs @@ -1332,6 +1332,56 @@ public async Task Dispose_ReleasesPendingGaugeInstrument() _ = await Assert.That(Volatile.Read(ref gaugeCompleted)).IsTrue(); } + [Test] + [NotInParallel("OutboxMetrics")] + public async Task Dispose_ReleasesCounterAndHistogramInstruments() + { + string[] names = + [ + "pulse.outbox.processed.total", + "pulse.outbox.failed.total", + "pulse.outbox.deadletter.total", + "pulse.outbox.processing.duration", + ]; + var published = new ConcurrentBag(); + var completed = new ConcurrentBag(); + using var meterListener = new MeterListener(); + meterListener.InstrumentPublished = (instrument, listener) => + { + if ( + string.Equals(instrument.Meter.Name, "NetEvolve.Pulse", StringComparison.Ordinal) + && names.Contains(instrument.Name, StringComparer.Ordinal) + ) + { + published.Add(instrument); + listener.EnableMeasurementEvents(instrument); + } + }; + meterListener.MeasurementsCompleted = (instrument, _) => completed.Add(instrument); + meterListener.Start(); + var publishedBefore = published.ToArray(); + + using var repository = new InMemoryOutboxRepository(); + var service = new OutboxProcessorHostedService( + CreateScopeFactory(repository), + new InMemoryMessageTransport(), + CreateLifetime(), + Options.Create(new OutboxProcessorOptions()), + CreateLogger(), + TimeProvider.System + ); + var created = published.Except(publishedBefore).ToArray(); + + service.Dispose(); + + // Other tests may create services in parallel, so look for one meter whose four instruments were released. + var released = created + .GroupBy(i => i.Meter) + .Any(g => g.Select(i => i.Name).Distinct().Count() == names.Length && g.All(completed.Contains)); + + _ = await Assert.That(released).IsTrue(); + } + [Test] public async Task ExecuteAsync_WithExponentialBackoffEnabled_SetsNextRetryAt(CancellationToken cancellationToken) { From 7425758630801786ff8937a8bb3b398b949d941c Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Martin=20St=C3=BChmer?= Date: Tue, 29 Sep 2026 02:12:22 +0200 Subject: [PATCH 11/14] fix(interceptors): record inner dispose and Current failures of a stream query as errors --- ...treamQueryInterceptor{TQuery,TResponse}.cs | 40 ++++++++++--------- 1 file changed, 22 insertions(+), 18 deletions(-) diff --git a/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsStreamQueryInterceptor{TQuery,TResponse}.cs b/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsStreamQueryInterceptor{TQuery,TResponse}.cs index 23ff5675..d114b1a6 100644 --- a/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsStreamQueryInterceptor{TQuery,TResponse}.cs +++ b/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsStreamQueryInterceptor{TQuery,TResponse}.cs @@ -167,10 +167,16 @@ public async IAsyncEnumerable HandleAsync( while (enumerator is not null) { - bool hasNext; + TResponse current; try { - hasNext = await enumerator.MoveNextAsync().ConfigureAwait(false); + if (!await enumerator.MoveNextAsync().ConfigureAwait(false)) + { + completed = true; + break; + } + + current = enumerator.Current; } catch (Exception ex) { @@ -178,33 +184,31 @@ public async IAsyncEnumerable HandleAsync( break; } - if (!hasNext) - { - completed = true; - break; - } - // yield return is valid here: it is inside try/finally but NOT inside try/catch - yield return enumerator.Current; + yield return current; } } finally { - try + if (enumerator is not null) { - if (enumerator is not null) + try { await enumerator.DisposeAsync().ConfigureAwait(false); } + catch (Exception ex) + { + // An earlier fault wins, so the thrown exception always matches the recorded error.type. + caughtExceptionInfo ??= ExceptionDispatchInfo.Capture(ex); + } } - finally - { - // Runs for every outcome, including a consumer that stops early and disposes the iterator. - RecordOutcome(activity, tags, startTime, caughtExceptionInfo?.SourceException, completed); - } - } - caughtExceptionInfo?.Throw(); + // Runs for every outcome, including a consumer that stops early and disposes the iterator. + RecordOutcome(activity, tags, startTime, caughtExceptionInfo?.SourceException, completed); + + // Thrown here and not after the finally block, because an early stop never reaches that code. + caughtExceptionInfo?.Throw(); + } } /// From 4ef491d4bcceab3c39e982cd31dee3e15e6c1506 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Martin=20St=C3=BChmer?= Date: Tue, 29 Sep 2026 02:17:07 +0200 Subject: [PATCH 12/14] fix(outbox): create the outbox counters and histogram on the instance meter so Dispose releases them --- .../Outbox/OutboxProcessorHostedService.cs | 18 ++++++++++-------- 1 file changed, 10 insertions(+), 8 deletions(-) diff --git a/src/NetEvolve.Pulse/Outbox/OutboxProcessorHostedService.cs b/src/NetEvolve.Pulse/Outbox/OutboxProcessorHostedService.cs index 2e2f76d1..89864e58 100644 --- a/src/NetEvolve.Pulse/Outbox/OutboxProcessorHostedService.cs +++ b/src/NetEvolve.Pulse/Outbox/OutboxProcessorHostedService.cs @@ -87,9 +87,10 @@ internal sealed partial class OutboxProcessorHostedService : BackgroundService private long _pendingCount; /// - /// Instance-scoped meter hosting the pulse.outbox.pending observable gauge. The gauge callback - /// captures this service instance, so the instrument lifetime is bound to the service lifetime by - /// disposing this meter in instead of publishing on the process-wide static meter. + /// Instance-scoped meter hosting the pulse.outbox.* counters, the duration histogram and the + /// pulse.outbox.pending observable gauge. The gauge callback captures this service instance, so the + /// instrument lifetime is bound to the service lifetime by disposing this meter in + /// instead of publishing on the process-wide static meter, which is never disposed. /// private readonly Meter _meter; @@ -130,29 +131,30 @@ public OutboxProcessorHostedService( _useSemanticConventionUnits = telemetryOptions?.Value.UseSemanticConventionUnits ?? false; var messageUnit = _useSemanticConventionUnits ? "{message}" : "messages"; - _processedCounter = Defaults.Meter.CreateCounter( + _meter = new Meter(Defaults.Meter.Name, Defaults.Version); + + _processedCounter = _meter.CreateCounter( "pulse.outbox.processed.total", messageUnit, "Cumulative number of successfully processed outbox messages." ); - _failedCounter = Defaults.Meter.CreateCounter( + _failedCounter = _meter.CreateCounter( "pulse.outbox.failed.total", messageUnit, "Cumulative number of failed outbox processing attempts." ); - _deadLetterCounter = Defaults.Meter.CreateCounter( + _deadLetterCounter = _meter.CreateCounter( "pulse.outbox.deadletter.total", messageUnit, "Cumulative number of outbox messages moved to dead-letter." ); _processingDurationHistogram = TelemetryUnits.CreateDurationHistogram( - Defaults.Meter, + _meter, "pulse.outbox.processing.duration", "each outbox processing batch", _useSemanticConventionUnits ); - _meter = new Meter(Defaults.Meter.Name, Defaults.Version); _ = _meter.CreateObservableGauge( "pulse.outbox.pending", observeValue: () => Volatile.Read(ref _pendingCount), From 14e64ee301d123c434a258236cdbe045f44d3f4f Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Martin=20St=C3=BChmer?= Date: Tue, 29 Sep 2026 02:17:07 +0200 Subject: [PATCH 13/14] fix(interceptors): share telemetry instruments per name and unit mode instead of creating them per instance --- ...ivityAndMetricsEventInterceptor{TEvent}.cs | 9 +-- ...sRequestInterceptor{TRequest,TResponse}.cs | 9 +-- ...treamQueryInterceptor{TQuery,TResponse}.cs | 9 +-- .../Internals/TelemetryUnits.cs | 68 +++++++++++++++++++ 4 files changed, 77 insertions(+), 18 deletions(-) diff --git a/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsEventInterceptor{TEvent}.cs b/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsEventInterceptor{TEvent}.cs index a0d202fb..6123b659 100644 --- a/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsEventInterceptor{TEvent}.cs +++ b/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsEventInterceptor{TEvent}.cs @@ -53,24 +53,21 @@ public ActivityAndMetricsEventInterceptor( { _timeProvider = timeProvider; _useSemanticConventionUnits = options?.Value.UseSemanticConventionUnits ?? false; - _eventCounter = TelemetryUnits.CreateCounter( - Defaults.Meter, + _eventCounter = TelemetryUnits.GetSharedCounter( "pulse.events.total", "events", "{event}", "Total number of events processed.", _useSemanticConventionUnits ); - _errorsCounter = TelemetryUnits.CreateCounter( - Defaults.Meter, + _errorsCounter = TelemetryUnits.GetSharedCounter( "pulse.event.errors", "errors", "{error}", "Total number of event errors.", _useSemanticConventionUnits ); - _eventDurationHistogram = TelemetryUnits.CreateDurationHistogram( - Defaults.Meter, + _eventDurationHistogram = TelemetryUnits.GetSharedDurationHistogram( "pulse.event.duration", "event processing", _useSemanticConventionUnits diff --git a/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsRequestInterceptor{TRequest,TResponse}.cs b/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsRequestInterceptor{TRequest,TResponse}.cs index 46112725..da627780 100644 --- a/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsRequestInterceptor{TRequest,TResponse}.cs +++ b/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsRequestInterceptor{TRequest,TResponse}.cs @@ -55,24 +55,21 @@ public ActivityAndMetricsRequestInterceptor( { _timeProvider = timeProvider; _useSemanticConventionUnits = options?.Value.UseSemanticConventionUnits ?? false; - _requestCounter = TelemetryUnits.CreateCounter( - Defaults.Meter, + _requestCounter = TelemetryUnits.GetSharedCounter( "pulse.requests.total", "requests", "{request}", "Total number of requests processed.", _useSemanticConventionUnits ); - _errorsCounter = TelemetryUnits.CreateCounter( - Defaults.Meter, + _errorsCounter = TelemetryUnits.GetSharedCounter( "pulse.request.errors", "errors", "{error}", "Total number of request errors.", _useSemanticConventionUnits ); - _requestDurationHistogram = TelemetryUnits.CreateDurationHistogram( - Defaults.Meter, + _requestDurationHistogram = TelemetryUnits.GetSharedDurationHistogram( "pulse.request.duration", "request processing", _useSemanticConventionUnits diff --git a/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsStreamQueryInterceptor{TQuery,TResponse}.cs b/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsStreamQueryInterceptor{TQuery,TResponse}.cs index d114b1a6..d42d4423 100644 --- a/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsStreamQueryInterceptor{TQuery,TResponse}.cs +++ b/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsStreamQueryInterceptor{TQuery,TResponse}.cs @@ -70,24 +70,21 @@ public ActivityAndMetricsStreamQueryInterceptor( { _timeProvider = timeProvider; _useSemanticConventionUnits = options?.Value.UseSemanticConventionUnits ?? false; - _streamQueryCounter = TelemetryUnits.CreateCounter( - Defaults.Meter, + _streamQueryCounter = TelemetryUnits.GetSharedCounter( "pulse.stream_query.total", "queries", "{query}", "Total number of stream queries processed.", _useSemanticConventionUnits ); - _errorsCounter = TelemetryUnits.CreateCounter( - Defaults.Meter, + _errorsCounter = TelemetryUnits.GetSharedCounter( "pulse.stream_query.errors", "errors", "{error}", "Total number of stream query errors.", _useSemanticConventionUnits ); - _streamQueryDurationHistogram = TelemetryUnits.CreateDurationHistogram( - Defaults.Meter, + _streamQueryDurationHistogram = TelemetryUnits.GetSharedDurationHistogram( "pulse.stream_query.duration", "stream query processing", _useSemanticConventionUnits diff --git a/src/NetEvolve.Pulse/Internals/TelemetryUnits.cs b/src/NetEvolve.Pulse/Internals/TelemetryUnits.cs index 0c47f341..831384e6 100644 --- a/src/NetEvolve.Pulse/Internals/TelemetryUnits.cs +++ b/src/NetEvolve.Pulse/Internals/TelemetryUnits.cs @@ -1,5 +1,6 @@ namespace NetEvolve.Pulse.Internals; +using System.Collections.Generic; using System.Diagnostics.Metrics; /// @@ -29,6 +30,54 @@ internal static class TelemetryUnits 10, ]; + /// + /// Instruments on , one per name and unit mode. The static meter is never disposed, + /// so creating instruments per interceptor instance or per service provider would accumulate them. + /// + private static readonly Dictionary<(string Name, bool UseSemanticConventionUnits), Instrument> SharedInstruments = + []; + + /// + /// Gets the counter on , creating it once per unit mode. + /// + /// The instrument name. + /// The legacy unit, for example requests. + /// The UCUM annotation unit, for example {request}. + /// The instrument description. + /// Whether to use the semantic convention unit. + /// The shared counter. + public static Counter GetSharedCounter( + string name, + string legacyUnit, + string annotationUnit, + string description, + bool useSemanticConventionUnits + ) => + GetShared( + name, + useSemanticConventionUnits, + () => + CreateCounter(Defaults.Meter, name, legacyUnit, annotationUnit, description, useSemanticConventionUnits) + ); + + /// + /// Gets the duration histogram on , creating it once per unit mode. + /// + /// The instrument name. + /// The measured operation used in the description, for example request processing. + /// Whether to record seconds. + /// The shared histogram. + public static Histogram GetSharedDurationHistogram( + string name, + string subject, + bool useSemanticConventionUnits + ) => + GetShared( + name, + useSemanticConventionUnits, + () => CreateDurationHistogram(Defaults.Meter, name, subject, useSemanticConventionUnits) + ); + /// /// Creates a counter on with the legacy or the UCUM annotation unit. /// @@ -86,4 +135,23 @@ bool useSemanticConventionUnits /// The duration in seconds or milliseconds. public static double ToDuration(TimeSpan elapsed, bool useSemanticConventionUnits) => useSemanticConventionUnits ? elapsed.TotalSeconds : elapsed.TotalMilliseconds; + + /// + /// Returns the cached instrument for and the unit mode, or creates and caches it. + /// The lock keeps concurrent first calls from creating the instrument twice. + /// + private static T GetShared(string name, bool useSemanticConventionUnits, Func create) + where T : Instrument + { + lock (SharedInstruments) + { + if (!SharedInstruments.TryGetValue((name, useSemanticConventionUnits), out var instrument)) + { + instrument = create(); + SharedInstruments.Add((name, useSemanticConventionUnits), instrument); + } + + return (T)instrument; + } + } } From 5d76936532a402c6376938d2f92a83229af0c621 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Martin=20St=C3=BChmer?= Date: Tue, 29 Sep 2026 02:17:41 +0200 Subject: [PATCH 14/14] docs(telemetry): describe cancelled and dispose-failed streams and fix the stale pulse.success remarks --- .../2026-09-28-opentelemetry-telemetry-conventions.md | 2 +- .../ActivityAndMetricsEventInterceptor{TEvent}.cs | 2 +- ...ityAndMetricsRequestInterceptor{TRequest,TResponse}.cs | 2 +- ...yAndMetricsStreamQueryInterceptor{TQuery,TResponse}.cs | 8 ++++++-- src/NetEvolve.Pulse/README.md | 3 ++- 5 files changed, 11 insertions(+), 6 deletions(-) diff --git a/decisions/2026-09-28-opentelemetry-telemetry-conventions.md b/decisions/2026-09-28-opentelemetry-telemetry-conventions.md index 11d84716..8d504a2d 100644 --- a/decisions/2026-09-28-opentelemetry-telemetry-conventions.md +++ b/decisions/2026-09-28-opentelemetry-telemetry-conventions.md @@ -42,7 +42,7 @@ No earlier decision defines the Pulse telemetry contract. Units are part of the * MUST leave the activity status `Unset` when an operation succeeds and when a stream consumer stops early. MUST set `Error` with the exception message on failure. * MUST set `error.type` to `Exception.GetType().FullName` on the span, the error counter and the duration histogram of a failed operation. The value MUST be identical on all three. The existing `pulse.exception.*` tags stay. -* MUST record `pulse.stream_query.duration` exactly once per stream query, from the `finally` block of the iterator. A stream the consumer abandons carries `pulse.stream.completed=false` on the histogram and the activity, keeps the status `Unset` and carries no `pulse.success` tag. Completed and faulted streams keep their `pulse.success` semantics. +* MUST record `pulse.stream_query.duration` exactly once per stream query, from the `finally` block of the iterator. A stream the consumer abandons carries `pulse.stream.completed=false` on the histogram and the activity, keeps the status `Unset` and carries no `pulse.success` tag. Completed and faulted streams keep their `pulse.success` semantics. A stream is faulted when its handler, `MoveNextAsync`, `Current` or the inner `DisposeAsync` throws, including an `OperationCanceledException` from a cancelled token. When the inner `DisposeAsync` throws after an earlier fault, the earlier exception is thrown and recorded, so the thrown exception and `error.type` always match. * MUST treat an exception thrown synchronously by the stream handler delegate like an exception thrown during enumeration. * MUST offer the new units through `ActivityAndMetricsOptions.UseSemanticConventionUnits` (default `false`). The option applies to the interceptors and to `OutboxProcessorHostedService`: diff --git a/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsEventInterceptor{TEvent}.cs b/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsEventInterceptor{TEvent}.cs index 6123b659..bd365433 100644 --- a/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsEventInterceptor{TEvent}.cs +++ b/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsEventInterceptor{TEvent}.cs @@ -83,7 +83,7 @@ public ActivityAndMetricsEventInterceptor( /// Increments event counter metrics /// Measures and records execution duration /// Captures exception details on failure - /// Marks success/failure status in both activity and metrics + /// Tags pulse.success on the activity and the duration histogram /// Leaves the activity status unless the handler fails /// Sets error.type on the activity, the error counter and the duration histogram on failure /// diff --git a/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsRequestInterceptor{TRequest,TResponse}.cs b/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsRequestInterceptor{TRequest,TResponse}.cs index da627780..cb3ec341 100644 --- a/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsRequestInterceptor{TRequest,TResponse}.cs +++ b/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsRequestInterceptor{TRequest,TResponse}.cs @@ -85,7 +85,7 @@ public ActivityAndMetricsRequestInterceptor( /// Increments request counter metrics /// Measures and records execution duration /// Captures exception details on failure - /// Marks success/failure status in both activity and metrics + /// Tags pulse.success on the activity and the duration histogram /// Leaves the activity status unless the handler fails /// Sets error.type on the activity, the error counter and the duration histogram on failure /// diff --git a/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsStreamQueryInterceptor{TQuery,TResponse}.cs b/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsStreamQueryInterceptor{TQuery,TResponse}.cs index d42d4423..5f3aaf80 100644 --- a/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsStreamQueryInterceptor{TQuery,TResponse}.cs +++ b/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsStreamQueryInterceptor{TQuery,TResponse}.cs @@ -101,14 +101,18 @@ public ActivityAndMetricsStreamQueryInterceptor( /// Increments stream query counter metrics /// Measures and records execution duration /// Captures exception details on failure - /// Marks success/failure status in both activity and metrics + /// Tags pulse.success on the activity and the duration histogram, except for abandoned streams /// Leaves the activity status unless the stream faults /// Sets error.type on the activity, the error counter and the duration histogram on failure /// /// Records the duration exactly once for every outcome. A stream whose consumer stops enumerating early - /// (for example , Take or a disconnected client) is tagged pulse.stream.completed=false + /// (for example or Take) is tagged pulse.stream.completed=false /// and carries no pulse.success tag /// + /// + /// Records an from a cancelled token as a failure, and treats an exception from + /// the inner enumerator's DisposeAsync or Current as a fault; an earlier fault wins over a dispose failure + /// /// Yields items unchanged without buffering /// /// diff --git a/src/NetEvolve.Pulse/README.md b/src/NetEvolve.Pulse/README.md index 1b965f7d..6f1f6f38 100644 --- a/src/NetEvolve.Pulse/README.md +++ b/src/NetEvolve.Pulse/README.md @@ -351,7 +351,8 @@ The custom serializer will be used for all payload operations within Pulse. Ensu * Successful operations leave the activity status `Unset`. Failed operations set `Error` with the exception message. * Failed operations carry `error.type` (the full exception type name) on the activity, the error counter and the duration histogram. -* A stream query records its duration once for every outcome. A stream whose consumer stops early (`break`, `Take`, a disconnected client) carries `pulse.stream.completed=false` instead of `pulse.success`. +* A stream query records its duration once for every outcome. A stream whose consumer stops early (`break`, `Take`) carries `pulse.stream.completed=false` instead of `pulse.success`. +* A stream query whose handler honours a cancelled token, for example `HttpContext.RequestAborted` after a client disconnects, fails with `OperationCanceledException` and is recorded as a failure with `error.type=System.OperationCanceledException`. An exception from the inner enumerator's `DisposeAsync` is also recorded as a failure, unless the stream had already faulted; the earlier exception is then both thrown and recorded. ### Semantic Convention Units