Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -110,7 +110,6 @@ CancellationToken ct

Logger.Current = Log;
var evt = DeserializeData(contentType, eventType, msg.Body, streamName);
var applicationProperties = msg.ApplicationProperties.Concat(MessageProperties(msg));

var ctx = new MessageConsumeContext(
msg.MessageId,
Expand All @@ -123,7 +122,7 @@ CancellationToken ct
run.NextSequence(),
msg.EnqueuedTime.UtcDateTime,
evt,
AsMeta(applicationProperties),
evt is null ? null : AsMeta(msg.ApplicationProperties.Concat(MessageProperties(msg))),
SubscriptionId,
ct
);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -245,9 +245,9 @@ RegisteredHandler<TAggregate, TState, TId> Build() {
);

Func<TCommand, IEventWriter> DefaultResolveWriter()
=> _ => Ensure.NotNull(writer, $"Function to resolve event writer from {typeof(TCommand).Name} is not defined and no default writer is set");
=> _ => writer ?? throw new ArgumentNullException($"Function to resolve event writer from {typeof(TCommand).Name} is not defined and no default writer is set");

Func<TCommand, IEventReader> DefaultResolveReader()
=> _ => Ensure.NotNull(reader, $"Function to resolve event reader from {typeof(TCommand).Name} is not defined and no default reader is set");
=> _ => reader ?? throw new ArgumentNullException($"Function to resolve event reader from {typeof(TCommand).Name} is not defined and no default reader is set");
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -18,10 +18,14 @@ class ApplicationEventSource : EventSource {
const int CommandHandlerRegisteredId = 5;

[NonEvent]
public void CommandHandlerNotFound<TCommand>() => CommandHandlerNotFound(typeof(TCommand).Name);
public void CommandHandlerNotFound<TCommand>() {
if (IsEnabled(EventLevel.Error, EventKeywords.All)) CommandHandlerNotFound(typeof(TCommand).Name);
}

[NonEvent]
public void ErrorHandlingCommand<TCommand>(Exception e) => ErrorHandlingCommand(typeof(TCommand).Name, e.ToString());
public void ErrorHandlingCommand<TCommand>(Exception e) {
if (IsEnabled(EventLevel.Error, EventKeywords.All)) ErrorHandlingCommand(typeof(TCommand).Name, e.ToString());
}

[NonEvent]
public void CommandHandled<TCommand>() {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -9,27 +9,33 @@ namespace Eventuous.Diagnostics;
using Tracing;

static class CommandServiceActivity {
// Shared by every traced service, whatever its state type: the metrics listener only observes the latest
// listener announced under a name, so a listener per service would hide the metrics of all but one.
static readonly DiagnosticSource MetricsSource = new DiagnosticListener(CommandServiceMetrics.ListenerName);

public static async Task<Result<T>> TryExecute<T, TCommand>(
string appServiceTypeName,
TCommand command,
DiagnosticSource diagnosticSource,
HandleCommand<T, TCommand> handleCommand,
CancellationToken cancellationToken
) where TCommand : class where T : State<T>, new() {
var cmdName = command.GetType().Name;

using var activity = StartActivity(appServiceTypeName, cmdName);
using var measure = Measure.Start(diagnosticSource, new CommandServiceMetricsContext(appServiceTypeName, cmdName));
// Nobody listening means nothing to record, so skip the measure and its context rather than allocate both.
using var measure = MetricsSource.IsEnabled(Measure.EventName)
? Measure.Start(MetricsSource, new CommandServiceMetricsContext(appServiceTypeName, cmdName))
: null;

try {
var result = await handleCommand(command, cancellationToken).NoContext();
activity?.SetActivityStatus(result is { Success: true } ? ActivityStatus.Ok() : ActivityStatus.Error(result.Exception));
if (!result.Success) measure.SetError();
if (!result.Success) measure?.SetError();

return result;
} catch (Exception e) {
activity?.SetActivityStatus(ActivityStatus.Error(e));
measure.SetError();
measure?.SetError();

throw;
}
Expand Down
Original file line number Diff line number Diff line change
@@ -1,17 +1,14 @@
// Copyright (C) Eventuous HQ OÜ. All rights reserved
// Licensed under the Apache License, Version 2.0.

using System.Diagnostics;

namespace Eventuous.Diagnostics;

public class TracedCommandService<TState> : ICommandService<TState> where TState : State<TState>, new() {
public static ICommandService<TState> Trace(ICommandService<TState> appService) => new TracedCommandService<TState>(appService);

ICommandService<TState> InnerService { get; }

readonly string _appServiceTypeName;
readonly DiagnosticSource _metricsSource = new DiagnosticListener(CommandServiceMetrics.ListenerName);
readonly string _appServiceTypeName;

TracedCommandService(ICommandService<TState> appService) {
_appServiceTypeName = appService.GetType().Name;
Expand All @@ -23,7 +20,6 @@ public Task<Result<TState>> Handle<TCommand>(TCommand command, CancellationToken
=> CommandServiceActivity.TryExecute(
_appServiceTypeName,
command,
_metricsSource,
InnerService.Handle,
cancellationToken
);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,6 @@ public async Task<Result<TState>> Handle<TCommand>(TCommand command, Cancellatio

result.ThrowIfError();

throw new ApplicationException($"Error handling command {command}");
return result;
}
}
3 changes: 2 additions & 1 deletion src/Core/src/Eventuous.Diagnostics/ActivityExtensions.cs
Original file line number Diff line number Diff line change
Expand Up @@ -5,9 +5,10 @@ namespace Eventuous.Diagnostics;

public static class ActivityExtensions {
extension(Activity activity) {
// Same result as searching Tags, which only yields string values, without allocating an enumerator and a closure.
[MethodImpl(MethodImplOptions.AggressiveInlining)]
string? GetParentTag(string tag)
=> activity.Parent?.Tags.FirstOrDefault(x => x.Key == tag).Value;
=> activity.Parent?.GetTagItem(tag) as string;

[MethodImpl(MethodImplOptions.AggressiveInlining)]
public Activity CopyParentTag(string tag, string? parentTag = null) {
Expand Down
5 changes: 4 additions & 1 deletion src/Core/src/Eventuous.Diagnostics/ActivityStatus.cs
Original file line number Diff line number Diff line change
Expand Up @@ -4,8 +4,11 @@
namespace Eventuous.Diagnostics;

public record ActivityStatus(ActivityStatusCode StatusCode, string? Description, Exception? Exception) {
static readonly ActivityStatus PlainOk = new(ActivityStatusCode.Ok, null, null);

// Shared when there's no description: the record is immutable, and this is the status of every successful span.
public static ActivityStatus Ok(string? description = null)
=> new(ActivityStatusCode.Ok, description, null);
=> description == null ? PlainOk : new(ActivityStatusCode.Ok, description, null);

public static ActivityStatus Error(Exception? exception = null, string? description = null)
=> new(ActivityStatusCode.Error, description ?? exception?.Message, exception);
Expand Down
5 changes: 2 additions & 3 deletions src/Core/src/Eventuous.Diagnostics/Metrics/Measure.cs
Original file line number Diff line number Diff line change
Expand Up @@ -16,12 +16,11 @@ public sealed class Measure(DiagnosticSource diagnosticSource, object context) :
Justification = "MeasureContext is not referencing anything."
)]
void Record() {
var stoppedAt = DateTime.UtcNow;
var duration = stoppedAt - _startedAt;
var duration = Stopwatch.GetElapsedTime(_startedAt);
diagnosticSource.Write(EventName, new MeasureContext(duration, _error, context));
}

readonly DateTime _startedAt = DateTime.UtcNow;
readonly long _startedAt = Stopwatch.GetTimestamp();

bool _error;

Expand Down
7 changes: 5 additions & 2 deletions src/Core/src/Eventuous.Domain/Aggregate.cs
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,7 @@ namespace Eventuous;
/// <summary>
/// Get the list of pending changes (new events) within the scope of the current operation.
/// </summary>
public IReadOnlyCollection<object> Changes => _changes.AsReadOnly();
public IReadOnlyCollection<object> Changes => _changesView ??= _changes.AsReadOnly();

/// <summary>
/// A collection with all the aggregate events, previously persisted and new
Expand All @@ -36,10 +36,13 @@ namespace Eventuous;
/// The current version is set to the original version when the aggregate is loaded from the store.
/// It should increase for each state transition performed within the scope of the current operation.
/// </summary>
public long CurrentVersion => OriginalVersion + Changes.Count;
public long CurrentVersion => OriginalVersion + _changes.Count;

readonly List<object> _changes = [];

// The view wraps the live list, so it can be created once and reused
IReadOnlyCollection<object>? _changesView;

/// <summary>
/// Adds an event to the list of pending changes.
/// </summary>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,7 @@ public abstract class BaseTracer {

protected async Task<T> Trace<T>(StreamName stream, string operation, Func<Task<T>> task) {
using var activity = StartActivity(stream, operation);
using var measure = Measure.Start(MetricsSource, new PersistenceMetricsContext(ComponentName, operation));
using var measure = StartMeasure(operation);

try {
var result = await task().NoContext();
Expand All @@ -27,22 +27,22 @@ protected async Task<T> Trace<T>(StreamName stream, string operation, Func<Task<
return result;
} catch (Exception e) {
activity?.SetActivityStatus(ActivityStatus.Error(e));
measure.SetError();
measure?.SetError();

throw;
}
}

protected async Task Trace(StreamName stream, string operation, Func<Task> task) {
using var activity = StartActivity(stream, operation);
using var measure = Measure.Start(MetricsSource, new PersistenceMetricsContext(ComponentName, operation));
using var measure = StartMeasure(operation);

try {
await task().NoContext();
activity?.SetActivityStatus(ActivityStatus.Ok());
} catch (Exception e) {
activity?.SetActivityStatus(ActivityStatus.Error(e));
measure.SetError();
measure?.SetError();

throw;
}
Expand All @@ -55,7 +55,7 @@ protected async IAsyncEnumerable<T> TraceEnumerable<T>(
[EnumeratorCancellation] CancellationToken cancellationToken = default
) {
using var activity = StartActivity(stream, operation);
using var measure = Measure.Start(MetricsSource, new PersistenceMetricsContext(ComponentName, operation));
using var measure = StartMeasure(operation);

var enumerator = source.GetAsyncEnumerator(cancellationToken);

Expand All @@ -67,7 +67,7 @@ protected async IAsyncEnumerable<T> TraceEnumerable<T>(
moved = await enumerator.MoveNextAsync().NoContext();
} catch (Exception e) {
activity?.SetActivityStatus(ActivityStatus.Error(e));
measure.SetError();
measure?.SetError();

throw;
}
Expand All @@ -81,6 +81,12 @@ protected async IAsyncEnumerable<T> TraceEnumerable<T>(
activity?.SetActivityStatus(ActivityStatus.Ok());
}

// Nobody listening means nothing to record, so skip the measure and its context rather than allocate both.
private protected Measure? StartMeasure(string operation)
=> MetricsSource.IsEnabled(Measure.EventName)
? Measure.Start(MetricsSource, new PersistenceMetricsContext(ComponentName, operation))
: null;

protected static Activity? StartActivity(StreamName stream, string operationName) {
if (!EventuousDiagnostics.Enabled) return null;

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,6 @@

namespace Eventuous.Diagnostics.Tracing;

using Metrics;
using static Constants;

public class TracedEventWriter(IEventWriter writer) : BaseTracer, IEventWriter {
Expand All @@ -21,8 +20,9 @@ CancellationToken cancellationToken
) {
using var activity = StartActivity(stream, Operations.AppendEvents);

using var measure = Measure.Start(MetricsSource, new PersistenceMetricsContext(ComponentName, Operations.AppendEvents));
using var measure = StartMeasure(Operations.AppendEvents);

// Snapshot, so the inner writer can't see the caller change the collection while the append is in flight
var tracedEvents = events
.Select(x => x with { Metadata = x.Metadata.AddActivityTags(activity) })
.ToArray();
Expand All @@ -34,7 +34,7 @@ CancellationToken cancellationToken
return result;
} catch (Exception e) {
activity?.SetActivityStatus(ActivityStatus.Error(e));
measure.SetError();
measure?.SetError();

throw;
}
Expand All @@ -46,7 +46,7 @@ public async Task<AppendEventsResult[]> AppendEvents(IReadOnlyCollection<NewStre
var streamNames = new StreamName(string.Join(", ", appends.Select(a => a.StreamName.ToString())));
using var activity = StartActivity(streamNames, Operations.AppendEvents);

using var measure = Measure.Start(MetricsSource, new PersistenceMetricsContext(ComponentName, Operations.AppendEvents));
using var measure = StartMeasure(Operations.AppendEvents);

var tracedAppends = appends.Select(a => a with { Events = [.. a.Events.Select(x => x with { Metadata = x.Metadata.AddActivityTags(activity) })] }).ToArray();

Expand All @@ -57,7 +57,7 @@ public async Task<AppendEventsResult[]> AppendEvents(IReadOnlyCollection<NewStre
return results;
} catch (Exception e) {
activity?.SetActivityStatus(ActivityStatus.Error(e));
measure.SetError();
measure?.SetError();

throw;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,11 @@ public class DefaultEventSerializer : IEventSerializer {
const string DynamicSerializationMessage =
"DefaultEventSerializer uses reflection-based System.Text.Json serialization. Use DefaultStaticEventSerializer with a JsonSerializerContext in trimmed or AOT applications.";

// Failure results are immutable, so one instance per error kind is shared instead of allocating per event
static readonly FailedToDeserialize UnknownType = new(DeserializationError.UnknownType);
static readonly FailedToDeserialize ContentTypeMismatch = new(DeserializationError.ContentTypeMismatch);
static readonly FailedToDeserialize PayloadEmpty = new(DeserializationError.PayloadEmpty);

readonly JsonSerializerOptions _options;
readonly ITypeMapper _typeMapper;

Expand All @@ -29,14 +34,14 @@ public DefaultEventSerializer(JsonSerializerOptions options, ITypeMapper? typeMa
public DeserializationResult DeserializeEvent(ReadOnlySpan<byte> data, string eventType, string contentType) {
var typeMapped = _typeMapper.TryGetType(eventType, out var dataType);

if (!typeMapped) return new FailedToDeserialize(DeserializationError.UnknownType);
if (contentType != ContentType) return new FailedToDeserialize(DeserializationError.ContentTypeMismatch);
if (!typeMapped) return UnknownType;
if (contentType != ContentType) return ContentTypeMismatch;

var deserialized = JsonSerializer.Deserialize(data, dataType!, _options);

return deserialized != null
? new SuccessfullyDeserialized(deserialized)
: new FailedToDeserialize(DeserializationError.PayloadEmpty);
: PayloadEmpty;
}

[UnconditionalSuppressMessage("Trimming", "IL2026", Justification = "The constructor is annotated with RequiresUnreferencedCode, so an instance only exists if the caller acknowledged the requirement")]
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -9,19 +9,24 @@ namespace Eventuous;

[PublicAPI]
public class DefaultStaticEventSerializer(JsonSerializerContext context, ITypeMapper? typeMapper = null) : IEventSerializer {
// Failure results are immutable, so one instance per error kind is shared instead of allocating per event
static readonly FailedToDeserialize UnknownType = new(DeserializationError.UnknownType);
static readonly FailedToDeserialize ContentTypeMismatch = new(DeserializationError.ContentTypeMismatch);
static readonly FailedToDeserialize PayloadEmpty = new(DeserializationError.PayloadEmpty);

readonly ITypeMapper _typeMapper = typeMapper ?? TypeMap.Instance;

public DeserializationResult DeserializeEvent(ReadOnlySpan<byte> data, string eventType, string contentType) {
var typeMapped = _typeMapper.TryGetType(eventType, out var dataType);

if (!typeMapped) return new FailedToDeserialize(DeserializationError.UnknownType);
if (contentType != ContentType) return new FailedToDeserialize(DeserializationError.ContentTypeMismatch);
if (!typeMapped) return UnknownType;
if (contentType != ContentType) return ContentTypeMismatch;

var deserialized = JsonSerializer.Deserialize(data, dataType!, context);

return deserialized != null
? new SuccessfullyDeserialized(deserialized)
: new FailedToDeserialize(DeserializationError.PayloadEmpty);
: PayloadEmpty;
}

public SerializationResult SerializeEvent(object evt)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -8,10 +8,10 @@ namespace Eventuous.Subscriptions.Consumers;
// ReSharper disable once ParameterTypeCanBeEnumerable.Local
public class DefaultConsumer(IEventHandler[] eventHandlers) : IMessageConsumer {
public async ValueTask Consume(IMessageConsumeContext context) {
var scope = new Dictionary<string, object> {
{"SubscriptionId", context.SubscriptionId},
{"Stream", context.Stream},
{"MessageType", context.MessageType}
var scope = new KeyValuePair<string, object>[] {
new("SubscriptionId", context.SubscriptionId),
new("Stream", context.Stream),
new("MessageType", context.MessageType)
};

using var _ = context.LogContext.Logger.BeginScope(scope);
Expand Down
Loading
Loading