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..8d504a2d --- /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/TelemetryUnits.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. 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`: + + | 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 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`. + +## 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`. 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 3c8a4cc7..bd365433 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,34 @@ 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.GetSharedCounter( + "pulse.events.total", + "events", + "{event}", + "Total number of events processed.", + _useSemanticConventionUnits + ); + _errorsCounter = TelemetryUnits.GetSharedCounter( + "pulse.event.errors", + "errors", + "{error}", + "Total number of event errors.", + _useSemanticConventionUnits + ); + _eventDurationHistogram = TelemetryUnits.GetSharedDurationHistogram( + "pulse.event.duration", + "event processing", + _useSemanticConventionUnits + ); + } /// /// @@ -62,7 +83,9 @@ internal sealed class ActivityAndMetricsEventInterceptor : IEventInterce /// 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 /// /// public async Task HandleAsync( @@ -93,7 +116,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 { @@ -102,33 +125,40 @@ 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); // 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) { 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( + TelemetryUnits.ToDuration(errorTime - startTime, _useSemanticConventionUnits), + [.. tags, new(Success, false), new(ErrorType, errorType)] + ); throw; } 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 942ba238..cb3ec341 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,34 @@ 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.GetSharedCounter( + "pulse.requests.total", + "requests", + "{request}", + "Total number of requests processed.", + _useSemanticConventionUnits + ); + _errorsCounter = TelemetryUnits.GetSharedCounter( + "pulse.request.errors", + "errors", + "{error}", + "Total number of request errors.", + _useSemanticConventionUnits + ); + _requestDurationHistogram = TelemetryUnits.GetSharedDurationHistogram( + "pulse.request.duration", + "request processing", + _useSemanticConventionUnits + ); + } /// /// @@ -64,7 +85,9 @@ internal sealed class 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 /// /// public async Task HandleAsync( @@ -102,7 +125,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 { @@ -111,35 +134,42 @@ 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); // 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; } 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( + TelemetryUnits.ToDuration(errorTime - startTime, _useSemanticConventionUnits), + [.. tags, new(Success, false), new(ErrorType, errorType)] + ); throw; } diff --git a/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsStreamQueryInterceptor{TQuery,TResponse}.cs b/src/NetEvolve.Pulse/Interceptors/ActivityAndMetricsStreamQueryInterceptor{TQuery,TResponse}.cs index d00fba0f..5f3aaf80 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,34 @@ 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.GetSharedCounter( + "pulse.stream_query.total", + "queries", + "{query}", + "Total number of stream queries processed.", + _useSemanticConventionUnits + ); + _errorsCounter = TelemetryUnits.GetSharedCounter( + "pulse.stream_query.errors", + "errors", + "{error}", + "Total number of stream query errors.", + _useSemanticConventionUnits + ); + _streamQueryDurationHistogram = TelemetryUnits.GetSharedDurationHistogram( + "pulse.stream_query.duration", + "stream query processing", + _useSemanticConventionUnits + ); + } /// /// @@ -80,7 +101,18 @@ internal sealed class ActivityAndMetricsStreamQueryInterceptorIncrements 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 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 /// /// @@ -115,21 +147,37 @@ 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 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; + TResponse current; try { - hasNext = await enumerator.MoveNextAsync().ConfigureAwait(false); + if (!await enumerator.MoveNextAsync().ConfigureAwait(false)) + { + completed = true; + break; + } + + current = enumerator.Current; } catch (Exception ex) { @@ -137,57 +185,83 @@ public async IAsyncEnumerable HandleAsync( break; } - if (!hasNext) - { - break; - } - // yield return is valid here: it is inside try/finally but NOT inside try/catch - yield return enumerator.Current; + yield return current; } } finally { - await enumerator.DisposeAsync().ConfigureAwait(false); + 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); + } + } + + // 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(); } + } + + /// + /// 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 = TelemetryUnits.ToDuration(endTime - startTime, _useSemanticConventionUnits); + _ = activity?.SetEndTime(endTime.UtcDateTime); - if (caughtExceptionInfo is not null) + 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"; } diff --git a/src/NetEvolve.Pulse/Internals/TelemetryUnits.cs b/src/NetEvolve.Pulse/Internals/TelemetryUnits.cs new file mode 100644 index 00000000..831384e6 --- /dev/null +++ b/src/NetEvolve.Pulse/Internals/TelemetryUnits.cs @@ -0,0 +1,157 @@ +namespace NetEvolve.Pulse.Internals; + +using System.Collections.Generic; +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, + ]; + + /// + /// 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. + /// + /// 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; + + /// + /// 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; + } + } +} diff --git a/src/NetEvolve.Pulse/Outbox/OutboxProcessorHostedService.cs b/src/NetEvolve.Pulse/Outbox/OutboxProcessorHostedService.cs index 78505688..6c2d05e6 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 @@ -98,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; @@ -113,13 +103,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 +128,37 @@ TimeProvider timeProvider _logger = logger; _timeProvider = timeProvider; + _useSemanticConventionUnits = telemetryOptions?.Value.UseSemanticConventionUnits ?? false; + var messageUnit = _useSemanticConventionUnits ? "{message}" : "messages"; + _meter = new Meter(Defaults.Meter.Name, Defaults.Version); + + _processedCounter = _meter.CreateCounter( + "pulse.outbox.processed.total", + messageUnit, + "Cumulative number of successfully processed outbox messages." + ); + _failedCounter = _meter.CreateCounter( + "pulse.outbox.failed.total", + messageUnit, + "Cumulative number of failed outbox processing attempts." + ); + _deadLetterCounter = _meter.CreateCounter( + "pulse.outbox.deadletter.total", + messageUnit, + "Cumulative number of outbox messages moved to dead-letter." + ); + _processingDurationHistogram = TelemetryUnits.CreateDurationHistogram( + _meter, + "pulse.outbox.processing.duration", + "each outbox processing batch", + _useSemanticConventionUnits + ); + _ = _meter.CreateObservableGauge( "pulse.outbox.pending", observeValue: () => Volatile.Read(ref _pendingCount), - unit: "messages", + unit: messageUnit, description: "Current number of pending outbox messages." ); } @@ -207,11 +225,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) { @@ -442,7 +463,7 @@ CancellationToken cancellationToken try { - ProcessedCounter.Add(1); + _processedCounter.Add(1); } catch (Exception ex) { @@ -466,7 +487,7 @@ CancellationToken cancellationToken try { - DeadLetterCounter.Add(1); + _deadLetterCounter.Add(1); } catch (Exception metricEx) { @@ -483,7 +504,7 @@ await repository try { - FailedCounter.Add(1); + _failedCounter.Add(1); } catch (Exception metricEx) { @@ -557,7 +578,7 @@ CancellationToken cancellationToken try { - ProcessedCounter.Add(messages.Length); + _processedCounter.Add(messages.Length); } catch (Exception ex) { @@ -628,12 +649,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 1814a2b2..7ffe05ce 100644 --- a/src/NetEvolve.Pulse/README.md +++ b/src/NetEvolve.Pulse/README.md @@ -367,6 +367,42 @@ 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`) 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 + +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. 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/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 d8e765a9..4b62d3c3 100644 --- a/tests/NetEvolve.Pulse.Tests.Unit/Interceptors/ActivityAndMetricsEventInterceptorTests.cs +++ b/tests/NetEvolve.Pulse.Tests.Unit/Interceptors/ActivityAndMetricsEventInterceptorTests.cs @@ -1,6 +1,10 @@ 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; using NetEvolve.Pulse.Extensibility; using NetEvolve.Pulse.Interceptors; @@ -54,7 +58,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 +81,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 +280,216 @@ 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(); + } + } + + [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"); + } + } + + [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(); + 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..96a58ceb 100644 --- a/tests/NetEvolve.Pulse.Tests.Unit/Interceptors/ActivityAndMetricsRequestInterceptorTests.cs +++ b/tests/NetEvolve.Pulse.Tests.Unit/Interceptors/ActivityAndMetricsRequestInterceptorTests.cs @@ -1,6 +1,10 @@ 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; using NetEvolve.Pulse.Extensibility; using NetEvolve.Pulse.Interceptors; @@ -126,7 +130,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 +155,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 +325,226 @@ 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(); + } + } + + [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"); + } + } + + [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; } + 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..142103f3 100644 --- a/tests/NetEvolve.Pulse.Tests.Unit/Interceptors/ActivityAndMetricsStreamQueryInterceptorTests.cs +++ b/tests/NetEvolve.Pulse.Tests.Unit/Interceptors/ActivityAndMetricsStreamQueryInterceptorTests.cs @@ -1,7 +1,11 @@ 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; using NetEvolve.Extensions.TUnit; using NetEvolve.Pulse.Extensibility; using NetEvolve.Pulse.Interceptors; @@ -51,7 +55,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 +85,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 +146,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 +176,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 +343,544 @@ 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"); + } + } + + [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"); + } + } + + [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); @@ -366,9 +908,46 @@ 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; } 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..f0d8b628 --- /dev/null +++ b/tests/NetEvolve.Pulse.Tests.Unit/Interceptors/PulseMeasurementCollector.cs @@ -0,0 +1,69 @@ +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. 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; + + public PulseMeasurementCollector(string? tagKey = null, string? tagValue = null) + { + _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 ( + _tagKey is null + || ( + 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 + ); +} diff --git a/tests/NetEvolve.Pulse.Tests.Unit/Outbox/OutboxProcessorHostedServiceTests.cs b/tests/NetEvolve.Pulse.Tests.Unit/Outbox/OutboxProcessorHostedServiceTests.cs index 83d12433..9aebd881 100644 --- a/tests/NetEvolve.Pulse.Tests.Unit/Outbox/OutboxProcessorHostedServiceTests.cs +++ b/tests/NetEvolve.Pulse.Tests.Unit/Outbox/OutboxProcessorHostedServiceTests.cs @@ -11,7 +11,9 @@ namespace NetEvolve.Pulse.Tests.Unit.Outbox; using NetEvolve.Extensions.TUnit; using NetEvolve.Pulse.Extensibility; using NetEvolve.Pulse.Extensibility.Outbox; +using NetEvolve.Pulse.Interceptors; using NetEvolve.Pulse.Outbox; +using NetEvolve.Pulse.Tests.Unit.Interceptors; using TUnit.Core; using TUnit.Mocks; @@ -1150,6 +1152,54 @@ public async Task ExecuteAsync_AfterProcessingCycle_RecordsProcessingDuration(Ca _ = await Assert.That(Volatile.Read(ref durationRecorded)).IsTrue(); } + [Test] + [NotInParallel("OutboxMetrics")] + public async Task ExecuteAsync_WithSemanticConventionUnits_RecordsSecondsAndMessageUnits( + CancellationToken cancellationToken + ) + { + cancellationToken.ThrowIfCancellationRequested(); + + using var collector = new PulseMeasurementCollector(); + + using var repository = new InMemoryOutboxRepository(); + await repository.AddAsync(CreateMessage(), cancellationToken).ConfigureAwait(false); + var transport = new InMemoryMessageTransport(); + var options = Options.Create(new OutboxProcessorOptions { PollingInterval = TimeSpan.FromMilliseconds(50) }); + using var service = new OutboxProcessorHostedService( + CreateScopeFactory(repository), + transport, + CreateLifetime(), + options, + CreateLogger(), + TimeProvider.System, + Options.Create(new ActivityAndMetricsOptions { UseSemanticConventionUnits = true }) + ); + + using var cts = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken); + await service.StartAsync(cts.Token).ConfigureAwait(false); + using var timeoutCts = CreateSignalTimeout(cancellationToken); + await repository.WaitForMarkingsAsync(1, 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); + + using (Assert.Multiple()) + { + _ = await Assert.That(collector.For("pulse.outbox.processing.duration").Any(m => 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( @@ -1282,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) {