Skip to content
Merged
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
14 changes: 14 additions & 0 deletions src/Aquila.Core/Abstractions/IDocumentStore.cs
Original file line number Diff line number Diff line change
Expand Up @@ -26,9 +26,23 @@ public interface IEventStore
void Append<TAggregate>(Guid streamId, long expectedVersion, params object[] events) where TAggregate : class;
void Append<TAggregate>(string streamId, long expectedVersion, params object[] events) where TAggregate : class;

void StartStreamTagged<TAggregate>(Guid streamId, IEnumerable<TaggedEvent> events) where TAggregate : class;
void StartStreamTagged<TAggregate>(string streamId, IEnumerable<TaggedEvent> events) where TAggregate : class;

void AppendTagged(Guid streamId, IEnumerable<TaggedEvent> events);
void AppendTagged(string streamId, IEnumerable<TaggedEvent> events);
void AppendTagged(Guid streamId, long expectedVersion, IEnumerable<TaggedEvent> events);
void AppendTagged(string streamId, long expectedVersion, IEnumerable<TaggedEvent> events);

void AppendTagged<TAggregate>(Guid streamId, IEnumerable<TaggedEvent> events) where TAggregate : class;
void AppendTagged<TAggregate>(string streamId, IEnumerable<TaggedEvent> events) where TAggregate : class;
void AppendTagged<TAggregate>(Guid streamId, long expectedVersion, IEnumerable<TaggedEvent> events) where TAggregate : class;
void AppendTagged<TAggregate>(string streamId, long expectedVersion, IEnumerable<TaggedEvent> events) where TAggregate : class;

Task<IReadOnlyList<IEvent>> FetchStreamAsync(Guid streamId, long fromVersion = 0, CancellationToken ct = default);
Task<IReadOnlyList<IEvent>> FetchStreamAsync(string streamId, long fromVersion = 0, CancellationToken ct = default);
Task<IReadOnlyList<IEvent>> FetchGlobalEventsAsync(long fromGlobalSequence, int batchSize = 1000, CancellationToken ct = default);
Task<IReadOnlyList<IEvent>> FetchEventsByTagAsync(string tag, long fromGlobalSequence = 0, int batchSize = 1000, CancellationToken ct = default);

Task<TAggregate?> AggregateStreamAsync<TAggregate>(Guid streamId, long version = 0, CancellationToken ct = default) where TAggregate : class, new();
Task<TAggregate?> AggregateStreamAsync<TAggregate>(string streamId, long version = 0, CancellationToken ct = default) where TAggregate : class, new();
Expand Down
10 changes: 10 additions & 0 deletions src/Aquila.Core/Events/IEvent.cs
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@ public interface IEvent
string? CorrelationId { get; set; }
string? CausationId { get; set; }
IReadOnlyDictionary<string, object> Headers { get; set; }
IReadOnlySet<string> Tags { get; set; }
}

/// <summary>
Expand All @@ -36,7 +37,10 @@ public interface IEvent<out T> : IEvent where T : class
/// </summary>
public sealed class EventEnvelope<T> : IEvent<T> where T : class
{
private static readonly IReadOnlySet<string> EmptyTags = new HashSet<string>();

private IReadOnlyDictionary<string, object>? _headers;
private IReadOnlySet<string>? _tags;

public Guid Id { get; set; } = Guid.NewGuid();
public string StreamId { get; set; } = string.Empty;
Expand All @@ -56,6 +60,12 @@ public IReadOnlyDictionary<string, object> Headers
get => _headers ?? ReadOnlyDictionary<string, object>.Empty;
set => _headers = value;
}

public IReadOnlySet<string> Tags
{
get => _tags ?? EmptyTags;
set => _tags = value;
}
}

/// <summary>
Expand Down
24 changes: 24 additions & 0 deletions src/Aquila.Core/Events/TaggedEvent.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,24 @@
namespace Aquila.Core.Events;

/// <summary>
/// Pairs a raw event payload with the tags it should carry when appended.
/// </summary>
public readonly struct TaggedEvent
{
private static readonly IReadOnlySet<string> EmptyTags = new HashSet<string>();

public object Data { get; }
public IReadOnlySet<string> Tags { get; }

public TaggedEvent(object data, IReadOnlySet<string>? tags = null)
{
ArgumentNullException.ThrowIfNull(data);
Data = data;
Tags = tags ?? EmptyTags;
}

public TaggedEvent(object data, IEnumerable<string> tags)
: this(data, tags is null ? EmptyTags : new HashSet<string>(tags))
{
}
}
2 changes: 2 additions & 0 deletions src/Aquila.Core/Events/UpcasterRegistry.cs
Original file line number Diff line number Diff line change
Expand Up @@ -127,6 +127,7 @@ private static IEvent CreateUpcastEnvelope(IEvent originalEvent, object newPaylo
var correlationIdProp = envelopeType.GetProperty(nameof(IEvent.CorrelationId))!;
var causationIdProp = envelopeType.GetProperty(nameof(IEvent.CausationId))!;
var headersProp = envelopeType.GetProperty(nameof(IEvent.Headers))!;
var tagsProp = envelopeType.GetProperty(nameof(IEvent.Tags))!;

var castPayload = Expression.Convert(payloadParam, t);
var eventTypeConst = Expression.Constant(t.FullName ?? t.Name);
Expand All @@ -146,6 +147,7 @@ private static IEvent CreateUpcastEnvelope(IEvent originalEvent, object newPaylo
Expression.Call(envVar, correlationIdProp.SetMethod!, Expression.Property(origParam, nameof(IEvent.CorrelationId))),
Expression.Call(envVar, causationIdProp.SetMethod!, Expression.Property(origParam, nameof(IEvent.CausationId))),
Expression.Call(envVar, headersProp.SetMethod!, Expression.Property(origParam, nameof(IEvent.Headers))),
Expression.Call(envVar, tagsProp.SetMethod!, Expression.Property(origParam, nameof(IEvent.Tags))),
Expression.Convert(envVar, typeof(IEvent))
);

Expand Down
156 changes: 136 additions & 20 deletions src/Aquila.Core/Sessions/QuerySession.cs
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,7 @@ public sealed class CoreEventStore : IEventStore
private readonly Dictionary<string, long> _streamExpectedVersions = new();
private readonly ConcurrentDictionary<string, Type> _streamAggregateTypes = new();

private static readonly ConcurrentDictionary<Type, Func<string, long, object, string, IEvent>> _envelopeFactories = new();
private static readonly ConcurrentDictionary<Type, Func<string, long, object, string, IReadOnlySet<string>, IEvent>> _envelopeFactories = new();
private static readonly ConcurrentDictionary<(Type AggregateType, Type EventType), Action<object, object>?> _applyMethodCache = new();

private readonly Func<(string? CorrelationId, string? CausationId, IReadOnlyDictionary<string, object> Headers)>? _headerProvider;
Expand Down Expand Up @@ -56,47 +56,86 @@ public CoreEventStore(
public void StartStream<TAggregate>(Guid streamId, params object[] events) where TAggregate : class
{
ArgumentNullException.ThrowIfNull(events);
StartStream<TAggregate>(streamId.ToString(), events);
StartStreamTagged<TAggregate>(streamId.ToString(), events.Select(e => new TaggedEvent(e)));
}

public void StartStream<TAggregate>(string streamId, params object[] events) where TAggregate : class
{
ArgumentException.ThrowIfNullOrWhiteSpace(streamId);
ArgumentNullException.ThrowIfNull(events);
StartStreamTagged<TAggregate>(streamId, events.Select(e => new TaggedEvent(e)));
}

public void StartStreamTagged<TAggregate>(Guid streamId, IEnumerable<TaggedEvent> events) where TAggregate : class
{
ArgumentNullException.ThrowIfNull(events);
StartStreamTagged<TAggregate>(streamId.ToString(), events);
}

public void StartStreamTagged<TAggregate>(string streamId, IEnumerable<TaggedEvent> events) where TAggregate : class
{
ArgumentException.ThrowIfNullOrWhiteSpace(streamId);
ArgumentNullException.ThrowIfNull(events);

_streamAggregateTypes[streamId] = typeof(TAggregate);

long version = 0;
foreach (var evt in events)
foreach (var tagged in events)
{
ArgumentNullException.ThrowIfNull(evt);
ArgumentNullException.ThrowIfNull(tagged.Data);
version++;
var envelope = CreateEnvelope(evt.GetType(), streamId, version, evt, _tenantId);
ApplyHeaders(envelope, evt);
var envelope = CreateEnvelope(tagged.Data.GetType(), streamId, version, tagged.Data, _tenantId, tagged.Tags);
ApplyHeaders(envelope, tagged.Data);
_uncommittedEvents.Add(envelope);
}
}

public void Append(Guid streamId, params object[] events)
{
ArgumentNullException.ThrowIfNull(events);
Append(streamId.ToString(), -1, events);
AppendTagged(streamId.ToString(), -1, events.Select(e => new TaggedEvent(e)));
}

public void Append(string streamId, params object[] events)
{
ArgumentException.ThrowIfNullOrWhiteSpace(streamId);
ArgumentNullException.ThrowIfNull(events);
Append(streamId, -1, events);
AppendTagged(streamId, -1, events.Select(e => new TaggedEvent(e)));
}

public void Append(Guid streamId, long expectedVersion, params object[] events)
{
ArgumentNullException.ThrowIfNull(events);
Append(streamId.ToString(), expectedVersion, events);
AppendTagged(streamId.ToString(), expectedVersion, events.Select(e => new TaggedEvent(e)));
}

public void Append(string streamId, long expectedVersion, params object[] events)
{
ArgumentException.ThrowIfNullOrWhiteSpace(streamId);
ArgumentNullException.ThrowIfNull(events);
AppendTagged(streamId, expectedVersion, events.Select(e => new TaggedEvent(e)));
}

public void AppendTagged(Guid streamId, IEnumerable<TaggedEvent> events)
{
ArgumentNullException.ThrowIfNull(events);
AppendTagged(streamId.ToString(), -1, events);
}

public void AppendTagged(string streamId, IEnumerable<TaggedEvent> events)
{
ArgumentException.ThrowIfNullOrWhiteSpace(streamId);
ArgumentNullException.ThrowIfNull(events);
AppendTagged(streamId, -1, events);
}

public void AppendTagged(Guid streamId, long expectedVersion, IEnumerable<TaggedEvent> events)
{
ArgumentNullException.ThrowIfNull(events);
AppendTagged(streamId.ToString(), expectedVersion, events);
}

public void AppendTagged(string streamId, long expectedVersion, IEnumerable<TaggedEvent> events)
{
ArgumentException.ThrowIfNullOrWhiteSpace(streamId);
ArgumentNullException.ThrowIfNull(events);
Expand All @@ -107,37 +146,93 @@ public void Append(string streamId, long expectedVersion, params object[] events
}

long version = expectedVersion > 0 ? expectedVersion : 0;
foreach (var evt in events)
foreach (var tagged in events)
{
ArgumentNullException.ThrowIfNull(evt);
ArgumentNullException.ThrowIfNull(tagged.Data);
version++;
var envelope = CreateEnvelope(evt.GetType(), streamId, version, evt, _tenantId);
ApplyHeaders(envelope, evt);
var envelope = CreateEnvelope(tagged.Data.GetType(), streamId, version, tagged.Data, _tenantId, tagged.Tags);
ApplyHeaders(envelope, tagged.Data);
_uncommittedEvents.Add(envelope);
}
}

public void Append<TAggregate>(Guid streamId, params object[] events) where TAggregate : class
{
ArgumentNullException.ThrowIfNull(events);
Append<TAggregate>(streamId.ToString(), -1, events);
AppendTagged<TAggregate>(streamId.ToString(), -1, events.Select(e => new TaggedEvent(e)));
}

public void Append<TAggregate>(string streamId, params object[] events) where TAggregate : class
{
ArgumentException.ThrowIfNullOrWhiteSpace(streamId);
ArgumentNullException.ThrowIfNull(events);
_streamAggregateTypes[streamId] = typeof(TAggregate);
Append(streamId, -1, events);
AppendTagged(streamId, -1, events.Select(e => new TaggedEvent(e)));
}

public void Append<TAggregate>(Guid streamId, long expectedVersion, params object[] events) where TAggregate : class
{
ArgumentNullException.ThrowIfNull(events);
Append<TAggregate>(streamId.ToString(), expectedVersion, events);
AppendTagged<TAggregate>(streamId.ToString(), expectedVersion, events.Select(e => new TaggedEvent(e)));
}

public void Append<TAggregate>(string streamId, long expectedVersion, params object[] events) where TAggregate : class
{
ArgumentException.ThrowIfNullOrWhiteSpace(streamId);
ArgumentNullException.ThrowIfNull(events);
_streamAggregateTypes[streamId] = typeof(TAggregate);
AppendTagged(streamId, expectedVersion, events.Select(e => new TaggedEvent(e)));
}

public void AppendTagged<TAggregate>(Guid streamId, IEnumerable<TaggedEvent> events) where TAggregate : class
{
ArgumentNullException.ThrowIfNull(events);
AppendTagged<TAggregate>(streamId.ToString(), -1, events);
}

public void AppendTagged<TAggregate>(string streamId, IEnumerable<TaggedEvent> events) where TAggregate : class
{
ArgumentException.ThrowIfNullOrWhiteSpace(streamId);
ArgumentNullException.ThrowIfNull(events);
_streamAggregateTypes[streamId] = typeof(TAggregate);
AppendTagged(streamId, -1, events);
}

public void AppendTagged<TAggregate>(Guid streamId, long expectedVersion, IEnumerable<TaggedEvent> events) where TAggregate : class
{
ArgumentNullException.ThrowIfNull(events);
AppendTagged<TAggregate>(streamId.ToString(), expectedVersion, events);
}

public void AppendTagged<TAggregate>(string streamId, long expectedVersion, IEnumerable<TaggedEvent> events) where TAggregate : class
{
ArgumentException.ThrowIfNullOrWhiteSpace(streamId);
ArgumentNullException.ThrowIfNull(events);
_streamAggregateTypes[streamId] = typeof(TAggregate);
AppendTagged(streamId, expectedVersion, events);
}

public void Append<TAggregate>(Guid streamId, IEnumerable<TaggedEvent> events) where TAggregate : class
{
ArgumentNullException.ThrowIfNull(events);
Append<TAggregate>(streamId.ToString(), -1, events);
}

public void Append<TAggregate>(string streamId, IEnumerable<TaggedEvent> events) where TAggregate : class
{
ArgumentException.ThrowIfNullOrWhiteSpace(streamId);
ArgumentNullException.ThrowIfNull(events);
_streamAggregateTypes[streamId] = typeof(TAggregate);
Append(streamId, -1, events);
}

public void Append<TAggregate>(Guid streamId, long expectedVersion, IEnumerable<TaggedEvent> events) where TAggregate : class
{
ArgumentNullException.ThrowIfNull(events);
Append<TAggregate>(streamId.ToString(), expectedVersion, events);
}

public void Append<TAggregate>(string streamId, long expectedVersion, IEnumerable<TaggedEvent> events) where TAggregate : class
{
ArgumentException.ThrowIfNullOrWhiteSpace(streamId);
ArgumentNullException.ThrowIfNull(events);
Expand Down Expand Up @@ -222,6 +317,22 @@ public async Task<IReadOnlyList<IEvent>> FetchGlobalEventsAsync(long fromGlobalS
return upcastEvents;
}

public async Task<IReadOnlyList<IEvent>> FetchEventsByTagAsync(string tag, long fromGlobalSequence = 0, int batchSize = 1000, CancellationToken ct = default)
{
var events = await _storage.FetchEventsByTagAsync(tag, fromGlobalSequence, batchSize, _tenantId, ct);
if (_upcasters == null || _upcasters.IsEmpty)
{
return events;
}

var upcastEvents = new List<IEvent>(events.Count);
foreach (var evt in events)
{
upcastEvents.Add(_upcasters.Upcast(evt));
}
return upcastEvents;
}

public Task<TAggregate?> AggregateStreamAsync<TAggregate>(Guid streamId, long version = 0, CancellationToken ct = default) where TAggregate : class, new()
{
return AggregateStreamAsync<TAggregate>(streamId.ToString(), version, ct);
Expand Down Expand Up @@ -269,14 +380,17 @@ public void ClearUncommittedEvents()
_streamExpectedVersions.Clear();
}

private static IEvent CreateEnvelope(Type eventType, string streamId, long version, object data, string tenantId)
private static readonly IReadOnlySet<string> EmptyTags = new HashSet<string>();

private static IEvent CreateEnvelope(Type eventType, string streamId, long version, object data, string tenantId, IReadOnlySet<string>? tags = null)
{
var factory = _envelopeFactories.GetOrAdd(eventType, t =>
{
var streamIdParam = Expression.Parameter(typeof(string), "streamId");
var versionParam = Expression.Parameter(typeof(long), "version");
var dataParam = Expression.Parameter(typeof(object), "data");
var tenantIdParam = Expression.Parameter(typeof(string), "tenantId");
var tagsParam = Expression.Parameter(typeof(IReadOnlySet<string>), "tags");

var envelopeType = typeof(EventEnvelope<>).MakeGenericType(t);
var ctor = Expression.New(envelopeType);
Expand All @@ -289,6 +403,7 @@ private static IEvent CreateEnvelope(Type eventType, string streamId, long versi
var eventTypeProp = envelopeType.GetProperty("EventType")!;
var dataProp = envelopeType.GetProperty("Data")!;
var tenantIdProp = envelopeType.GetProperty("TenantId")!;
var tagsProp = envelopeType.GetProperty("Tags")!;

var newGuidCall = Expression.Call(typeof(Guid), nameof(Guid.NewGuid), Type.EmptyTypes);
var castData = Expression.Convert(dataParam, t);
Expand All @@ -304,14 +419,15 @@ private static IEvent CreateEnvelope(Type eventType, string streamId, long versi
Expression.Call(envelopeVar, eventTypeProp.SetMethod!, eventTypeConst),
Expression.Call(envelopeVar, dataProp.SetMethod!, castData),
Expression.Call(envelopeVar, tenantIdProp.SetMethod!, tenantIdParam),
Expression.Call(envelopeVar, tagsProp.SetMethod!, tagsParam),
Expression.Convert(envelopeVar, typeof(IEvent))
);

return Expression.Lambda<Func<string, long, object, string, IEvent>>(
block, streamIdParam, versionParam, dataParam, tenantIdParam).Compile();
return Expression.Lambda<Func<string, long, object, string, IReadOnlySet<string>, IEvent>>(
block, streamIdParam, versionParam, dataParam, tenantIdParam, tagsParam).Compile();
});

return factory(streamId, version, data, tenantId);
return factory(streamId, version, data, tenantId, tags ?? EmptyTags);
}

internal static void ApplyEventToAggregate(object aggregate, object eventData)
Expand Down
Loading
Loading