diff --git a/src/Aquila.Core/Abstractions/IDocumentStore.cs b/src/Aquila.Core/Abstractions/IDocumentStore.cs index 197e278..97e5daa 100644 --- a/src/Aquila.Core/Abstractions/IDocumentStore.cs +++ b/src/Aquila.Core/Abstractions/IDocumentStore.cs @@ -26,9 +26,23 @@ public interface IEventStore void Append(Guid streamId, long expectedVersion, params object[] events) where TAggregate : class; void Append(string streamId, long expectedVersion, params object[] events) where TAggregate : class; + void StartStreamTagged(Guid streamId, IEnumerable events) where TAggregate : class; + void StartStreamTagged(string streamId, IEnumerable events) where TAggregate : class; + + void AppendTagged(Guid streamId, IEnumerable events); + void AppendTagged(string streamId, IEnumerable events); + void AppendTagged(Guid streamId, long expectedVersion, IEnumerable events); + void AppendTagged(string streamId, long expectedVersion, IEnumerable events); + + void AppendTagged(Guid streamId, IEnumerable events) where TAggregate : class; + void AppendTagged(string streamId, IEnumerable events) where TAggregate : class; + void AppendTagged(Guid streamId, long expectedVersion, IEnumerable events) where TAggregate : class; + void AppendTagged(string streamId, long expectedVersion, IEnumerable events) where TAggregate : class; + Task> FetchStreamAsync(Guid streamId, long fromVersion = 0, CancellationToken ct = default); Task> FetchStreamAsync(string streamId, long fromVersion = 0, CancellationToken ct = default); Task> FetchGlobalEventsAsync(long fromGlobalSequence, int batchSize = 1000, CancellationToken ct = default); + Task> FetchEventsByTagAsync(string tag, long fromGlobalSequence = 0, int batchSize = 1000, CancellationToken ct = default); Task AggregateStreamAsync(Guid streamId, long version = 0, CancellationToken ct = default) where TAggregate : class, new(); Task AggregateStreamAsync(string streamId, long version = 0, CancellationToken ct = default) where TAggregate : class, new(); diff --git a/src/Aquila.Core/Events/IEvent.cs b/src/Aquila.Core/Events/IEvent.cs index 3a81894..a8c9c9b 100644 --- a/src/Aquila.Core/Events/IEvent.cs +++ b/src/Aquila.Core/Events/IEvent.cs @@ -21,6 +21,7 @@ public interface IEvent string? CorrelationId { get; set; } string? CausationId { get; set; } IReadOnlyDictionary Headers { get; set; } + IReadOnlySet Tags { get; set; } } /// @@ -36,7 +37,10 @@ public interface IEvent : IEvent where T : class /// public sealed class EventEnvelope : IEvent where T : class { + private static readonly IReadOnlySet EmptyTags = new HashSet(); + private IReadOnlyDictionary? _headers; + private IReadOnlySet? _tags; public Guid Id { get; set; } = Guid.NewGuid(); public string StreamId { get; set; } = string.Empty; @@ -56,6 +60,12 @@ public IReadOnlyDictionary Headers get => _headers ?? ReadOnlyDictionary.Empty; set => _headers = value; } + + public IReadOnlySet Tags + { + get => _tags ?? EmptyTags; + set => _tags = value; + } } /// diff --git a/src/Aquila.Core/Events/TaggedEvent.cs b/src/Aquila.Core/Events/TaggedEvent.cs new file mode 100644 index 0000000..d16ad40 --- /dev/null +++ b/src/Aquila.Core/Events/TaggedEvent.cs @@ -0,0 +1,24 @@ +namespace Aquila.Core.Events; + +/// +/// Pairs a raw event payload with the tags it should carry when appended. +/// +public readonly struct TaggedEvent +{ + private static readonly IReadOnlySet EmptyTags = new HashSet(); + + public object Data { get; } + public IReadOnlySet Tags { get; } + + public TaggedEvent(object data, IReadOnlySet? tags = null) + { + ArgumentNullException.ThrowIfNull(data); + Data = data; + Tags = tags ?? EmptyTags; + } + + public TaggedEvent(object data, IEnumerable tags) + : this(data, tags is null ? EmptyTags : new HashSet(tags)) + { + } +} diff --git a/src/Aquila.Core/Events/UpcasterRegistry.cs b/src/Aquila.Core/Events/UpcasterRegistry.cs index 28a4116..97f31a5 100644 --- a/src/Aquila.Core/Events/UpcasterRegistry.cs +++ b/src/Aquila.Core/Events/UpcasterRegistry.cs @@ -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); @@ -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)) ); diff --git a/src/Aquila.Core/Sessions/QuerySession.cs b/src/Aquila.Core/Sessions/QuerySession.cs index 8d47758..fa14771 100644 --- a/src/Aquila.Core/Sessions/QuerySession.cs +++ b/src/Aquila.Core/Sessions/QuerySession.cs @@ -20,7 +20,7 @@ public sealed class CoreEventStore : IEventStore private readonly Dictionary _streamExpectedVersions = new(); private readonly ConcurrentDictionary _streamAggregateTypes = new(); - private static readonly ConcurrentDictionary> _envelopeFactories = new(); + private static readonly ConcurrentDictionary, IEvent>> _envelopeFactories = new(); private static readonly ConcurrentDictionary<(Type AggregateType, Type EventType), Action?> _applyMethodCache = new(); private readonly Func<(string? CorrelationId, string? CausationId, IReadOnlyDictionary Headers)>? _headerProvider; @@ -56,10 +56,23 @@ public CoreEventStore( public void StartStream(Guid streamId, params object[] events) where TAggregate : class { ArgumentNullException.ThrowIfNull(events); - StartStream(streamId.ToString(), events); + StartStreamTagged(streamId.ToString(), events.Select(e => new TaggedEvent(e))); } public void StartStream(string streamId, params object[] events) where TAggregate : class + { + ArgumentException.ThrowIfNullOrWhiteSpace(streamId); + ArgumentNullException.ThrowIfNull(events); + StartStreamTagged(streamId, events.Select(e => new TaggedEvent(e))); + } + + public void StartStreamTagged(Guid streamId, IEnumerable events) where TAggregate : class + { + ArgumentNullException.ThrowIfNull(events); + StartStreamTagged(streamId.ToString(), events); + } + + public void StartStreamTagged(string streamId, IEnumerable events) where TAggregate : class { ArgumentException.ThrowIfNullOrWhiteSpace(streamId); ArgumentNullException.ThrowIfNull(events); @@ -67,12 +80,12 @@ public void StartStream(string streamId, params object[] events) whe _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); } } @@ -80,23 +93,49 @@ public void StartStream(string streamId, params object[] events) whe 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 events) + { + ArgumentNullException.ThrowIfNull(events); + AppendTagged(streamId.ToString(), -1, events); + } + + public void AppendTagged(string streamId, IEnumerable events) + { + ArgumentException.ThrowIfNullOrWhiteSpace(streamId); + ArgumentNullException.ThrowIfNull(events); + AppendTagged(streamId, -1, events); + } + + public void AppendTagged(Guid streamId, long expectedVersion, IEnumerable events) + { + ArgumentNullException.ThrowIfNull(events); + AppendTagged(streamId.ToString(), expectedVersion, events); + } + + public void AppendTagged(string streamId, long expectedVersion, IEnumerable events) { ArgumentException.ThrowIfNullOrWhiteSpace(streamId); ArgumentNullException.ThrowIfNull(events); @@ -107,12 +146,12 @@ 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); } } @@ -120,7 +159,7 @@ public void Append(string streamId, long expectedVersion, params object[] events public void Append(Guid streamId, params object[] events) where TAggregate : class { 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) where TAggregate : class @@ -128,16 +167,72 @@ public void Append(string streamId, params object[] events) where TA 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(Guid streamId, long expectedVersion, params object[] events) where TAggregate : class { 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) 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(Guid streamId, IEnumerable events) where TAggregate : class + { + ArgumentNullException.ThrowIfNull(events); + AppendTagged(streamId.ToString(), -1, events); + } + + public void AppendTagged(string streamId, IEnumerable events) where TAggregate : class + { + ArgumentException.ThrowIfNullOrWhiteSpace(streamId); + ArgumentNullException.ThrowIfNull(events); + _streamAggregateTypes[streamId] = typeof(TAggregate); + AppendTagged(streamId, -1, events); + } + + public void AppendTagged(Guid streamId, long expectedVersion, IEnumerable events) where TAggregate : class + { + ArgumentNullException.ThrowIfNull(events); + AppendTagged(streamId.ToString(), expectedVersion, events); + } + + public void AppendTagged(string streamId, long expectedVersion, IEnumerable events) where TAggregate : class + { + ArgumentException.ThrowIfNullOrWhiteSpace(streamId); + ArgumentNullException.ThrowIfNull(events); + _streamAggregateTypes[streamId] = typeof(TAggregate); + AppendTagged(streamId, expectedVersion, events); + } + + public void Append(Guid streamId, IEnumerable events) where TAggregate : class + { + ArgumentNullException.ThrowIfNull(events); + Append(streamId.ToString(), -1, events); + } + + public void Append(string streamId, IEnumerable events) where TAggregate : class + { + ArgumentException.ThrowIfNullOrWhiteSpace(streamId); + ArgumentNullException.ThrowIfNull(events); + _streamAggregateTypes[streamId] = typeof(TAggregate); + Append(streamId, -1, events); + } + + public void Append(Guid streamId, long expectedVersion, IEnumerable events) where TAggregate : class + { + ArgumentNullException.ThrowIfNull(events); + Append(streamId.ToString(), expectedVersion, events); + } + + public void Append(string streamId, long expectedVersion, IEnumerable events) where TAggregate : class { ArgumentException.ThrowIfNullOrWhiteSpace(streamId); ArgumentNullException.ThrowIfNull(events); @@ -222,6 +317,22 @@ public async Task> FetchGlobalEventsAsync(long fromGlobalS return upcastEvents; } + public async Task> 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(events.Count); + foreach (var evt in events) + { + upcastEvents.Add(_upcasters.Upcast(evt)); + } + return upcastEvents; + } + public Task AggregateStreamAsync(Guid streamId, long version = 0, CancellationToken ct = default) where TAggregate : class, new() { return AggregateStreamAsync(streamId.ToString(), version, ct); @@ -269,7 +380,9 @@ public void ClearUncommittedEvents() _streamExpectedVersions.Clear(); } - private static IEvent CreateEnvelope(Type eventType, string streamId, long version, object data, string tenantId) + private static readonly IReadOnlySet EmptyTags = new HashSet(); + + private static IEvent CreateEnvelope(Type eventType, string streamId, long version, object data, string tenantId, IReadOnlySet? tags = null) { var factory = _envelopeFactories.GetOrAdd(eventType, t => { @@ -277,6 +390,7 @@ private static IEvent CreateEnvelope(Type eventType, string streamId, long versi 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), "tags"); var envelopeType = typeof(EventEnvelope<>).MakeGenericType(t); var ctor = Expression.New(envelopeType); @@ -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); @@ -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>( - block, streamIdParam, versionParam, dataParam, tenantIdParam).Compile(); + return Expression.Lambda, 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) diff --git a/src/Aquila.Core/Storage/InMemoryEventStorageProvider.cs b/src/Aquila.Core/Storage/InMemoryEventStorageProvider.cs index 6e60462..22225cb 100644 --- a/src/Aquila.Core/Storage/InMemoryEventStorageProvider.cs +++ b/src/Aquila.Core/Storage/InMemoryEventStorageProvider.cs @@ -105,6 +105,34 @@ public Task> FetchGlobalEventsAsync(long fromGlobalSequenc return Task.FromResult>(allEvents); } + public Task> FetchEventsByTagAsync(string tag, long fromGlobalSequence = 0, int batchSize = 1000, string? tenantId = null, CancellationToken ct = default) + { + ArgumentException.ThrowIfNullOrWhiteSpace(tag); + + if (batchSize <= 0) + { + return Task.FromResult>(Array.Empty()); + } + + List allEvents; + lock (_eventStreams) + { + allEvents = _eventStreams.Values + .SelectMany(s => + { + lock (s) { return s.ToList(); } + }) + .Where(e => (string.IsNullOrEmpty(tenantId) || e.TenantId == tenantId) + && e.GlobalSequence > fromGlobalSequence + && e.Tags.Contains(tag)) + .OrderBy(e => e.GlobalSequence) + .Take(batchSize) + .ToList(); + } + + return Task.FromResult>(allEvents); + } + public Task GetStreamHeaderAsync(string streamId, string? tenantId = null, CancellationToken ct = default) { ArgumentException.ThrowIfNullOrWhiteSpace(streamId); diff --git a/src/Aquila.Core/Storage/InMemoryStorageProvider.cs b/src/Aquila.Core/Storage/InMemoryStorageProvider.cs index 943cd86..462e863 100644 --- a/src/Aquila.Core/Storage/InMemoryStorageProvider.cs +++ b/src/Aquila.Core/Storage/InMemoryStorageProvider.cs @@ -74,6 +74,9 @@ public Task> FetchEventsAsync(string streamId, string? ten public Task> FetchGlobalEventsAsync(long fromGlobalSequence, int batchSize = 1000, string? tenantId = null, CancellationToken ct = default) => _events.FetchGlobalEventsAsync(fromGlobalSequence, batchSize, tenantId, ct); + public Task> FetchEventsByTagAsync(string tag, long fromGlobalSequence = 0, int batchSize = 1000, string? tenantId = null, CancellationToken ct = default) => + _events.FetchEventsByTagAsync(tag, fromGlobalSequence, batchSize, tenantId, ct); + public Task GetStreamHeaderAsync(string streamId, string? tenantId = null, CancellationToken ct = default) => _events.GetStreamHeaderAsync(streamId, tenantId, ct); diff --git a/src/Aquila.Core/Storage/StorageContracts.cs b/src/Aquila.Core/Storage/StorageContracts.cs index f69ce0f..5402679 100644 --- a/src/Aquila.Core/Storage/StorageContracts.cs +++ b/src/Aquila.Core/Storage/StorageContracts.cs @@ -170,6 +170,7 @@ public interface IEventStorageProvider : IDisposable, IAsyncDisposable Task AppendEventsAsync(string streamId, IEnumerable events, long expectedVersion, CancellationToken ct = default); Task> FetchEventsAsync(string streamId, string? tenantId = null, long fromVersion = 0, CancellationToken ct = default); Task> FetchGlobalEventsAsync(long fromGlobalSequence, int batchSize = 1000, string? tenantId = null, CancellationToken ct = default); + Task> FetchEventsByTagAsync(string tag, long fromGlobalSequence = 0, int batchSize = 1000, string? tenantId = null, CancellationToken ct = default); Task GetStreamHeaderAsync(string streamId, string? tenantId = null, CancellationToken ct = default); Task SaveSnapshotAsync(string streamId, long version, TAggregate snapshot, string tenantId = "default", CancellationToken ct = default) where TAggregate : class; Task<(TAggregate? Snapshot, long SnapshotVersion)> GetSnapshotAsync(string streamId, string tenantId = "default", CancellationToken ct = default) where TAggregate : class; diff --git a/src/Aquila.Cosmos/Events/CosmosEventStore.cs b/src/Aquila.Cosmos/Events/CosmosEventStore.cs index 7537445..ad93d90 100644 --- a/src/Aquila.Cosmos/Events/CosmosEventStore.cs +++ b/src/Aquila.Cosmos/Events/CosmosEventStore.cs @@ -56,6 +56,36 @@ public void Append(Guid streamId, long expectedVersion, params objec public void Append(string streamId, long expectedVersion, params object[] events) where TAggregate : class => _innerStore.Append(streamId, expectedVersion, events); + public void StartStreamTagged(Guid streamId, IEnumerable events) where TAggregate : class => + _innerStore.StartStreamTagged(streamId, events); + + public void StartStreamTagged(string streamId, IEnumerable events) where TAggregate : class => + _innerStore.StartStreamTagged(streamId, events); + + public void AppendTagged(Guid streamId, IEnumerable events) => + _innerStore.AppendTagged(streamId, events); + + public void AppendTagged(string streamId, IEnumerable events) => + _innerStore.AppendTagged(streamId, events); + + public void AppendTagged(Guid streamId, long expectedVersion, IEnumerable events) => + _innerStore.AppendTagged(streamId, expectedVersion, events); + + public void AppendTagged(string streamId, long expectedVersion, IEnumerable events) => + _innerStore.AppendTagged(streamId, expectedVersion, events); + + public void AppendTagged(Guid streamId, IEnumerable events) where TAggregate : class => + _innerStore.AppendTagged(streamId, events); + + public void AppendTagged(string streamId, IEnumerable events) where TAggregate : class => + _innerStore.AppendTagged(streamId, events); + + public void AppendTagged(Guid streamId, long expectedVersion, IEnumerable events) where TAggregate : class => + _innerStore.AppendTagged(streamId, expectedVersion, events); + + public void AppendTagged(string streamId, long expectedVersion, IEnumerable events) where TAggregate : class => + _innerStore.AppendTagged(streamId, expectedVersion, events); + public Task> FetchStreamAsync(Guid streamId, long fromVersion = 0, CancellationToken ct = default) => _innerStore.FetchStreamAsync(streamId, fromVersion, ct); @@ -65,6 +95,9 @@ public Task> FetchStreamAsync(string streamId, long fromVe public Task> FetchGlobalEventsAsync(long fromGlobalSequence, int batchSize = 1000, CancellationToken ct = default) => _innerStore.FetchGlobalEventsAsync(fromGlobalSequence, batchSize, ct); + public Task> FetchEventsByTagAsync(string tag, long fromGlobalSequence = 0, int batchSize = 1000, CancellationToken ct = default) => + _innerStore.FetchEventsByTagAsync(tag, fromGlobalSequence, batchSize, ct); + public Task AggregateStreamAsync(Guid streamId, long version = 0, CancellationToken ct = default) where TAggregate : class, new() => _innerStore.AggregateStreamAsync(streamId, version, ct); diff --git a/src/Aquila.Cosmos/Storage/CosmosEventStorageProvider.cs b/src/Aquila.Cosmos/Storage/CosmosEventStorageProvider.cs index a46573e..1e2917c 100644 --- a/src/Aquila.Cosmos/Storage/CosmosEventStorageProvider.cs +++ b/src/Aquila.Cosmos/Storage/CosmosEventStorageProvider.cs @@ -470,6 +470,159 @@ private async Task> FetchGlobalEventsUnsortedFallbackAsync } } + public async Task> FetchEventsByTagAsync(string tag, long fromGlobalSequence = 0, int batchSize = 1000, string? tenantId = null, CancellationToken ct = default) + { + ArgumentException.ThrowIfNullOrWhiteSpace(tag); + if (batchSize <= 0) + { + return Array.Empty(); + } + + var queryText = string.IsNullOrEmpty(tenantId) + ? "SELECT * FROM c WHERE c._docType = '$event' AND ARRAY_CONTAINS(c.data.Tags, @tag) AND c.data.GlobalSequence > @fromGlobalSequence ORDER BY c.data.GlobalSequence" + : "SELECT * FROM c WHERE c._docType = '$event' AND c._tenantId = @tenantId AND ARRAY_CONTAINS(c.data.Tags, @tag) AND c.data.GlobalSequence > @fromGlobalSequence ORDER BY c.data.GlobalSequence"; + + var queryDef = new QueryDefinition(queryText) + .WithParameter("@tag", tag) + .WithParameter("@fromGlobalSequence", fromGlobalSequence); + + if (!string.IsNullOrEmpty(tenantId)) + { + queryDef = queryDef.WithParameter("@tenantId", tenantId); + } + + var requestOptions = new QueryRequestOptions { MaxItemCount = batchSize }; + + var events = new List(); + using var iterator = EventContainer.GetItemQueryIterator>(queryDef, requestOptions: requestOptions); + if (iterator == null) return events; + + double globalCharge = 0.0; + try + { + while (iterator.HasMoreResults && events.Count < batchSize) + { + var response = await iterator.ReadNextAsync(ct).ConfigureAwait(false); + globalCharge += response.RequestCharge; + foreach (var item in response) + { + IEvent? @event = item.Data as IEvent; + if (@event == null && item.Data != null) + { + var rawJson = item.Data.ToString(); + if (!string.IsNullOrEmpty(rawJson)) + { + var envelope = Newtonsoft.Json.JsonConvert.DeserializeObject>(rawJson, PrivateConstructorContractResolver.Settings); + if (envelope != null) + { + _eventTypeResolver.EnsureTypedPayload(envelope); + } + @event = envelope; + } + } + + if (@event != null) + { + if (long.TryParse(item.Version, out var itemVer) && itemVer > 0 && @event.Version == 0) + { + @event.SetVersion(itemVer); + } + + if ((string.IsNullOrEmpty(tenantId) || item.TenantId == tenantId || @event.TenantId == tenantId) + && @event.GlobalSequence > fromGlobalSequence + && @event.Tags.Contains(tag)) + { + events.Add(@event); + if (events.Count >= batchSize) break; + } + } + } + } + + RecordCharge(globalCharge); + return events.OrderBy(e => e.GlobalSequence).Take(batchSize).ToList(); + } + catch (CosmosException ex) when (ex.StatusCode == System.Net.HttpStatusCode.InternalServerError || ex.StatusCode == System.Net.HttpStatusCode.BadRequest) + { + _logger?.LogWarning(ex, "Server-side sorted FetchEventsByTagAsync query failed ({StatusCode}). Falling back to unsorted server-side query with client sort.", ex.StatusCode); + return await FetchEventsByTagUnsortedFallbackAsync(tag, fromGlobalSequence, batchSize, tenantId, ct).ConfigureAwait(false); + } + } + + private async Task> FetchEventsByTagUnsortedFallbackAsync(string tag, long fromGlobalSequence, int batchSize, string? tenantId, CancellationToken ct) + { + var sql = "SELECT * FROM c WHERE c._docType = '$event' AND ARRAY_CONTAINS(c.data.Tags, @tag) AND c.data.GlobalSequence > @fromGlobalSequence"; + if (!string.IsNullOrEmpty(tenantId)) + { + sql += " AND c._tenantId = @tenantId"; + } + + var queryDef = new QueryDefinition(sql) + .WithParameter("@tag", tag) + .WithParameter("@fromGlobalSequence", fromGlobalSequence); + + if (!string.IsNullOrEmpty(tenantId)) + { + queryDef = queryDef.WithParameter("@tenantId", tenantId); + } + + var requestOptions = new QueryRequestOptions { MaxItemCount = batchSize }; + var events = new List(); + + using var iterator = EventContainer.GetItemQueryIterator>(queryDef, requestOptions: requestOptions); + if (iterator == null) return events; + + double totalCharge = 0.0; + try + { + while (iterator.HasMoreResults) + { + var response = await iterator.ReadNextAsync(ct).ConfigureAwait(false); + totalCharge += response.RequestCharge; + foreach (var item in response) + { + IEvent? @event = item.Data as IEvent; + if (@event == null && item.Data != null) + { + var rawJson = item.Data.ToString(); + if (!string.IsNullOrEmpty(rawJson)) + { + var envelope = Newtonsoft.Json.JsonConvert.DeserializeObject>(rawJson, PrivateConstructorContractResolver.Settings); + if (envelope != null) + { + _eventTypeResolver.EnsureTypedPayload(envelope); + } + @event = envelope; + } + } + + if (@event != null) + { + if (long.TryParse(item.Version, out var itemVer) && itemVer > 0 && @event.Version == 0) + { + @event.SetVersion(itemVer); + } + + if ((string.IsNullOrEmpty(tenantId) || item.TenantId == tenantId || @event.TenantId == tenantId) + && @event.GlobalSequence > fromGlobalSequence + && @event.Tags.Contains(tag)) + { + events.Add(@event); + } + } + } + } + + RecordCharge(totalCharge); + return events.OrderBy(e => e.GlobalSequence).Take(batchSize).ToList(); + } + catch (Exception ex) + { + _logger?.LogError(ex, "Failed to fetch events by tag with unsorted fallback query."); + return events.OrderBy(e => e.GlobalSequence).Take(batchSize).ToList(); + } + } + private async Task> FetchGlobalEventsDocTypeFallbackAsync(long fromGlobalSequence, int batchSize, string? tenantId, CancellationToken ct) { var sql = "SELECT * FROM c WHERE c._docType = '$event'"; diff --git a/src/Aquila.Cosmos/Storage/CosmosStorageProvider.cs b/src/Aquila.Cosmos/Storage/CosmosStorageProvider.cs index 37107f7..1a8f82d 100644 --- a/src/Aquila.Cosmos/Storage/CosmosStorageProvider.cs +++ b/src/Aquila.Cosmos/Storage/CosmosStorageProvider.cs @@ -132,6 +132,7 @@ public static ContainerProperties CreateDefaultEventsContainerProperties(string props.IndexingPolicy.IncludedPaths.Add(new IncludedPath { Path = "/_docType/?" }); props.IndexingPolicy.IncludedPaths.Add(new IncludedPath { Path = "/_tenantId/?" }); props.IndexingPolicy.IncludedPaths.Add(new IncludedPath { Path = "/data/GlobalSequence/?" }); + props.IndexingPolicy.IncludedPaths.Add(new IncludedPath { Path = "/data/Tags/*" }); props.IndexingPolicy.IncludedPaths.Add(new IncludedPath { Path = "/pk/?" }); props.IndexingPolicy.IncludedPaths.Add(new IncludedPath { Path = "/id/?" }); @@ -208,6 +209,9 @@ public Task> FetchEventsAsync(string streamId, string? ten public Task> FetchGlobalEventsAsync(long fromGlobalSequence, int batchSize = 1000, string? tenantId = null, CancellationToken ct = default) => _events.FetchGlobalEventsAsync(fromGlobalSequence, batchSize, tenantId, ct); + public Task> FetchEventsByTagAsync(string tag, long fromGlobalSequence = 0, int batchSize = 1000, string? tenantId = null, CancellationToken ct = default) => + _events.FetchEventsByTagAsync(tag, fromGlobalSequence, batchSize, tenantId, ct); + public Task GetStreamHeaderAsync(string streamId, string? tenantId = null, CancellationToken ct = default) => _events.GetStreamHeaderAsync(streamId, tenantId, ct); diff --git a/tests/Aquila.Core.Tests/Events/AppendTaggedTests.cs b/tests/Aquila.Core.Tests/Events/AppendTaggedTests.cs new file mode 100644 index 0000000..14f3a99 --- /dev/null +++ b/tests/Aquila.Core.Tests/Events/AppendTaggedTests.cs @@ -0,0 +1,47 @@ +using Shouldly; +using Aquila.Core.Configuration; +using Aquila.Core.Events; +using Aquila.Core.Sessions; +using Aquila.Core.Storage; + +namespace Aquila.Core.Tests; + +public sealed class AppendTaggedTests +{ + [Fact] + public async Task AppendTagged_ProducesEnvelopesWithSuppliedTags() + { + var storage = new InMemoryStorageProvider(); + var options = new StoreOptions { DocumentStorage = storage, EventStorage = storage }; + using var session = new DocumentSession(storage, storage, options); + + var streamId = Guid.NewGuid(); + var createdEvent = new AccountCreatedEvent(streamId, "Alice", 500m); + + session.Events.AppendTagged(streamId, new[] { new TaggedEvent(createdEvent, new HashSet { "AccountCreated", "Account" }) }); + await session.SaveChangesAsync(TestContext.Current.CancellationToken); + + var events = await storage.FetchEventsAsync(streamId.ToString(), ct: TestContext.Current.CancellationToken); + events.Count.ShouldBe(1); + events[0].Tags.ShouldContain("AccountCreated"); + events[0].Tags.ShouldContain("Account"); + } + + [Fact] + public async Task Append_ParamsObjectOverload_StillProducesEmptyTags() + { + var storage = new InMemoryStorageProvider(); + var options = new StoreOptions { DocumentStorage = storage, EventStorage = storage }; + using var session = new DocumentSession(storage, storage, options); + + var streamId = Guid.NewGuid(); + var createdEvent = new AccountCreatedEvent(streamId, "Bob", 100m); + + session.Events.Append(streamId, createdEvent); + await session.SaveChangesAsync(TestContext.Current.CancellationToken); + + var events = await storage.FetchEventsAsync(streamId.ToString(), ct: TestContext.Current.CancellationToken); + events.Count.ShouldBe(1); + events[0].Tags.ShouldBeEmpty(); + } +} diff --git a/tests/Aquila.Core.Tests/Events/EventExtensionsTests.cs b/tests/Aquila.Core.Tests/Events/EventExtensionsTests.cs index db00599..ab01ac8 100644 --- a/tests/Aquila.Core.Tests/Events/EventExtensionsTests.cs +++ b/tests/Aquila.Core.Tests/Events/EventExtensionsTests.cs @@ -20,6 +20,7 @@ private class ReadOnlyGlobalSequenceEvent : IEvent public string? CorrelationId { get; set; } public string? CausationId { get; set; } public IReadOnlyDictionary Headers { get; set; } = ReadOnlyDictionary.Empty; + public IReadOnlySet Tags { get; set; } = new HashSet(); } private class SamplePayload { } diff --git a/tests/Aquila.Core.Tests/Storage/InMemoryEventTagStorageTests.cs b/tests/Aquila.Core.Tests/Storage/InMemoryEventTagStorageTests.cs new file mode 100644 index 0000000..f4df9e9 --- /dev/null +++ b/tests/Aquila.Core.Tests/Storage/InMemoryEventTagStorageTests.cs @@ -0,0 +1,139 @@ +using Shouldly; +using Aquila.Core.Events; +using Aquila.Core.Storage; + +namespace Aquila.Core.Tests; + +public sealed class InMemoryEventTagStorageTests +{ + private static EventEnvelope CreateTaggedEvent(string streamId, long version, IReadOnlySet? tags = null, string tenantId = "default") => + new() + { + StreamId = streamId, + Version = version, + TenantId = tenantId, + Data = new AccountCreatedEvent(Guid.NewGuid(), "Alice", 100m), + Tags = tags ?? new HashSet() + }; + + [Fact] + public async Task FetchEventsByTagAsync_ReturnsOnlyEventsWithMatchingTag() + { + var provider = new InMemoryStorageProvider(); + var streamId = "stream-tags-1"; + + var patientCreated = CreateTaggedEvent(streamId, 1, new HashSet { "PatientCreated", "Patient" }); + var patientUpdated = CreateTaggedEvent(streamId, 2, new HashSet { "Patient" }); + var untagged = CreateTaggedEvent(streamId, 3); + + await provider.AppendEventsAsync(streamId, new IEvent[] { patientCreated, patientUpdated, untagged }, 0, TestContext.Current.CancellationToken); + + var createdOnly = await provider.FetchEventsByTagAsync("PatientCreated", ct: TestContext.Current.CancellationToken); + createdOnly.Count.ShouldBe(1); + createdOnly[0].Version.ShouldBe(1); + + var allPatient = await provider.FetchEventsByTagAsync("Patient", ct: TestContext.Current.CancellationToken); + allPatient.Count.ShouldBe(2); + } + + [Fact] + public async Task FetchEventsByTagAsync_UntaggedLegacyEvent_NeverMatchesAnyTag_AndDoesNotThrow() + { + var provider = new InMemoryStorageProvider(); + var streamId = "stream-tags-legacy"; + + var legacyEvent = new EventEnvelope + { + StreamId = streamId, + Version = 1, + Data = new AccountCreatedEvent(Guid.NewGuid(), "Alice", 100m) + }; + + legacyEvent.Tags.ShouldBeEmpty(); + + await provider.AppendEventsAsync(streamId, new[] { legacyEvent }, 0, TestContext.Current.CancellationToken); + + var result = await provider.FetchEventsByTagAsync("AnyTag", ct: TestContext.Current.CancellationToken); + result.ShouldBeEmpty(); + } + + [Fact] + public async Task FetchEventsByTagAsync_RespectsTenantIdFilter() + { + var provider = new InMemoryStorageProvider(); + + var tenantAEvent = CreateTaggedEvent("stream-a", 1, new HashSet { "Patient" }, tenantId: "tenant-a"); + var tenantBEvent = CreateTaggedEvent("stream-b", 1, new HashSet { "Patient" }, tenantId: "tenant-b"); + + await provider.AppendEventsAsync("stream-a", new IEvent[] { tenantAEvent }, 0, TestContext.Current.CancellationToken); + await provider.AppendEventsAsync("stream-b", new IEvent[] { tenantBEvent }, 0, TestContext.Current.CancellationToken); + + var tenantAResults = await provider.FetchEventsByTagAsync("Patient", tenantId: "tenant-a", ct: TestContext.Current.CancellationToken); + tenantAResults.Count.ShouldBe(1); + tenantAResults[0].TenantId.ShouldBe("tenant-a"); + } + + [Fact] + public async Task FetchEventsByTagAsync_RespectsFromGlobalSequenceExclusiveLowerBound() + { + var provider = new InMemoryStorageProvider(); + var streamId = "stream-tags-seq"; + + var first = CreateTaggedEvent(streamId, 1, new HashSet { "Patient" }); + var second = CreateTaggedEvent(streamId, 2, new HashSet { "Patient" }); + + await provider.AppendEventsAsync(streamId, new IEvent[] { first, second }, 0, TestContext.Current.CancellationToken); + + var fromFirst = await provider.FetchEventsByTagAsync("Patient", fromGlobalSequence: first.GlobalSequence, ct: TestContext.Current.CancellationToken); + fromFirst.Count.ShouldBe(1); + fromFirst[0].GlobalSequence.ShouldBe(second.GlobalSequence); + } + + [Fact] + public async Task FetchEventsByTagAsync_RespectsBatchSize() + { + var provider = new InMemoryStorageProvider(); + var streamId = "stream-tags-batch"; + + var events = Enumerable.Range(1, 5) + .Select(i => (IEvent)CreateTaggedEvent(streamId, i, new HashSet { "Patient" })) + .ToArray(); + + await provider.AppendEventsAsync(streamId, events, 0, TestContext.Current.CancellationToken); + + var page = await provider.FetchEventsByTagAsync("Patient", batchSize: 2, ct: TestContext.Current.CancellationToken); + page.Count.ShouldBe(2); + page[0].Version.ShouldBe(1); + page[1].Version.ShouldBe(2); + } + + [Fact] + public async Task FetchEventsByTagAsync_BatchSizeZeroOrNegative_ReturnsEmpty() + { + var provider = new InMemoryStorageProvider(); + + var result = await provider.FetchEventsByTagAsync("Patient", batchSize: 0, ct: TestContext.Current.CancellationToken); + result.ShouldBeEmpty(); + } + + [Fact] + public async Task FetchEventsByTagAsync_UnknownTag_ReturnsEmpty() + { + var provider = new InMemoryStorageProvider(); + var streamId = "stream-tags-unknown"; + + var tagged = CreateTaggedEvent(streamId, 1, new HashSet { "Patient" }); + await provider.AppendEventsAsync(streamId, new IEvent[] { tagged }, 0, TestContext.Current.CancellationToken); + + var result = await provider.FetchEventsByTagAsync("NoSuchTag", ct: TestContext.Current.CancellationToken); + result.ShouldBeEmpty(); + } + + [Fact] + public async Task FetchEventsByTagAsync_NullOrWhitespaceTag_Throws() + { + var provider = new InMemoryStorageProvider(); + + await Should.ThrowAsync(() => provider.FetchEventsByTagAsync(" ")); + } +} diff --git a/tests/Aquila.Core.Tests/Storage/StorageContractsTests.cs b/tests/Aquila.Core.Tests/Storage/StorageContractsTests.cs index 264a185..c5ea71b 100644 --- a/tests/Aquila.Core.Tests/Storage/StorageContractsTests.cs +++ b/tests/Aquila.Core.Tests/Storage/StorageContractsTests.cs @@ -373,6 +373,7 @@ private sealed class MinimalEventStorageProviderStub : IEventStorageProvider public Task AppendEventsAsync(string streamId, IEnumerable events, long expectedVersion, CancellationToken ct = default) => Task.CompletedTask; public Task> FetchEventsAsync(string streamId, string? tenantId = null, long fromVersion = 0, CancellationToken ct = default) => Task.FromResult>(Array.Empty()); public Task> FetchGlobalEventsAsync(long fromGlobalSequence, int batchSize = 1000, string? tenantId = null, CancellationToken ct = default) => Task.FromResult>(Array.Empty()); + public Task> FetchEventsByTagAsync(string tag, long fromGlobalSequence = 0, int batchSize = 1000, string? tenantId = null, CancellationToken ct = default) => Task.FromResult>(Array.Empty()); public Task GetStreamHeaderAsync(string streamId, string? tenantId = null, CancellationToken ct = default) => Task.FromResult(null); public Task SaveSnapshotAsync(string streamId, long version, TAggregate snapshot, string tenantId = "default", CancellationToken ct = default) where TAggregate : class => Task.CompletedTask; public Task<(TAggregate? Snapshot, long SnapshotVersion)> GetSnapshotAsync(string streamId, string tenantId = "default", CancellationToken ct = default) where TAggregate : class => Task.FromResult<(TAggregate?, long)>((null, 0)); diff --git a/tests/Aquila.Cosmos.Tests/Projections/CosmosDaemonTests.cs b/tests/Aquila.Cosmos.Tests/Projections/CosmosDaemonTests.cs index 6f808a2..9373334 100644 --- a/tests/Aquila.Cosmos.Tests/Projections/CosmosDaemonTests.cs +++ b/tests/Aquila.Cosmos.Tests/Projections/CosmosDaemonTests.cs @@ -105,6 +105,7 @@ public sealed class TestReadOnlyEvent : IEvent public string? CorrelationId { get; set; } public string? CausationId { get; set; } public IReadOnlyDictionary Headers { get; set; } = new Dictionary(); + public IReadOnlySet Tags { get; set; } = new HashSet(); public object Data { get; } public TestReadOnlyEvent(object data) diff --git a/tests/Aquila.Cosmos.Tests/Storage/CosmosEventStorageProviderTests.cs b/tests/Aquila.Cosmos.Tests/Storage/CosmosEventStorageProviderTests.cs index ead483a..c636f13 100644 --- a/tests/Aquila.Cosmos.Tests/Storage/CosmosEventStorageProviderTests.cs +++ b/tests/Aquila.Cosmos.Tests/Storage/CosmosEventStorageProviderTests.cs @@ -264,6 +264,164 @@ public async Task FetchGlobalEventsAsync_QueriesGlobalEvents_WithPagination() _provider.LastRequestCharge.ShouldBe(4.0); } + // ========================================== + // 4b. FetchEventsByTagAsync Tests + // ========================================== + + [Fact] + public async Task FetchEventsByTagAsync_BatchSizeZeroOrNegative_ReturnsEmpty() + { + var resultZero = await _provider.FetchEventsByTagAsync("Patient", batchSize: 0, ct: TestContext.Current.CancellationToken); + resultZero.ShouldBeEmpty(); + + var resultNeg = await _provider.FetchEventsByTagAsync("Patient", batchSize: -5, ct: TestContext.Current.CancellationToken); + resultNeg.ShouldBeEmpty(); + } + + [Fact] + public async Task FetchEventsByTagAsync_NullOrWhitespaceTag_Throws() + { + await Should.ThrowAsync(() => _provider.FetchEventsByTagAsync(" ", ct: TestContext.Current.CancellationToken)); + } + + [Fact] + public async Task FetchEventsByTagAsync_BuildsParameterizedQuery_WithArrayContainsAndTagParameter() + { + var iterator = Substitute.For>>(); + var page = Substitute.For>>(); + page.GetEnumerator().Returns(new List>().GetEnumerator()); + page.RequestCharge.Returns(1.0); + + iterator.HasMoreResults.Returns(true, false); + iterator.ReadNextAsync(Arg.Any()).Returns(Task.FromResult(page)); + + QueryDefinition? captured = null; + _mockEventContainer.GetItemQueryIterator>( + Arg.Do(q => captured = q), + requestOptions: Arg.Any()) + .Returns(iterator); + + await _provider.FetchEventsByTagAsync("Patient", ct: TestContext.Current.CancellationToken); + + captured.ShouldNotBeNull(); + captured!.QueryText.ShouldContain("ARRAY_CONTAINS(c.data.Tags, @tag)"); + captured.QueryText.ShouldNotContain("'Patient'"); + } + + [Fact] + public async Task FetchEventsByTagAsync_WithTenantId_IncludesTenantFilterInQueryText() + { + var iterator = Substitute.For>>(); + var page = Substitute.For>>(); + page.GetEnumerator().Returns(new List>().GetEnumerator()); + page.RequestCharge.Returns(1.0); + + iterator.HasMoreResults.Returns(true, false); + iterator.ReadNextAsync(Arg.Any()).Returns(Task.FromResult(page)); + + QueryDefinition? captured = null; + _mockEventContainer.GetItemQueryIterator>( + Arg.Do(q => captured = q), + requestOptions: Arg.Any()) + .Returns(iterator); + + await _provider.FetchEventsByTagAsync("Patient", tenantId: "t1", ct: TestContext.Current.CancellationToken); + + captured.ShouldNotBeNull(); + captured!.QueryText.ShouldContain("c._tenantId = @tenantId"); + } + + [Fact] + public async Task FetchEventsByTagAsync_QueriesAndReturnsMatchingEvents_WithPagination() + { + var evt1 = new EventEnvelope { StreamId = "s1", Version = 1, TenantId = "t1", GlobalSequence = 10, Tags = new HashSet { "Patient" }, Data = new TestEventPayload("o1", 10m) }; + var evt2 = new EventEnvelope { StreamId = "s2", Version = 1, TenantId = "t1", GlobalSequence = 20, Tags = new HashSet { "Patient" }, Data = new TestEventPayload("o2", 20m) }; + + var env1 = new CosmosDocumentEnvelope { Id = "e1", PartitionKey = "s1", TenantId = "t1", Data = evt1 }; + var env2 = new CosmosDocumentEnvelope { Id = "e2", PartitionKey = "s2", TenantId = "t1", Data = evt2 }; + + var iterator = Substitute.For>>(); + var page = Substitute.For>>(); + page.GetEnumerator().Returns(new List> { env1, env2 }.GetEnumerator()); + page.RequestCharge.Returns(4.0); + + iterator.HasMoreResults.Returns(true, false); + iterator.ReadNextAsync(Arg.Any()).Returns(Task.FromResult(page)); + + _mockEventContainer.GetItemQueryIterator>( + Arg.Any(), + requestOptions: Arg.Any()) + .Returns(iterator); + + var events = await _provider.FetchEventsByTagAsync("Patient", fromGlobalSequence: 5, batchSize: 10, tenantId: "t1", ct: TestContext.Current.CancellationToken); + + events.Count.ShouldBe(2); + events[0].GlobalSequence.ShouldBe(10); + events[1].GlobalSequence.ShouldBe(20); + _provider.LastRequestCharge.ShouldBe(4.0); + } + + [Fact] + public async Task FetchEventsByTagAsync_FallbacksToUnsortedQuery_WhenSortedQueryFails() + { + var evt = new EventEnvelope { StreamId = "s-fb", Version = 1, TenantId = "t1", GlobalSequence = 55, Tags = new HashSet { "Patient" }, Data = new TestEventPayload("o-fb", 10m) }; + var env = new CosmosDocumentEnvelope { Id = "e-fb", PartitionKey = "s-fb", TenantId = "t1", Data = evt }; + + var badRequestEx = new CosmosException("Bad Request", HttpStatusCode.BadRequest, 0, "act-1", 0); + + var sortedIterator = Substitute.For>>(); + sortedIterator.HasMoreResults.Returns(true); + sortedIterator.ReadNextAsync(Arg.Any()).Returns(Task.FromException>>(badRequestEx)); + + var fallbackIterator = Substitute.For>>(); + var page = Substitute.For>>(); + page.GetEnumerator().Returns(new List> { env }.GetEnumerator()); + fallbackIterator.HasMoreResults.Returns(true, false); + fallbackIterator.ReadNextAsync(Arg.Any()).Returns(Task.FromResult(page)); + + _mockEventContainer.GetItemQueryIterator>( + Arg.Is(q => q.QueryText.Contains("ORDER BY")), + requestOptions: Arg.Any()) + .Returns(sortedIterator); + + _mockEventContainer.GetItemQueryIterator>( + Arg.Is(q => !q.QueryText.Contains("ORDER BY")), + requestOptions: Arg.Any()) + .Returns(fallbackIterator); + + var events = await _provider.FetchEventsByTagAsync("Patient", fromGlobalSequence: 0, batchSize: 10, tenantId: "t1", ct: TestContext.Current.CancellationToken); + events.Count.ShouldBe(1); + events[0].GlobalSequence.ShouldBe(55); + } + + [Fact] + public async Task FetchEventsByTagAsync_RawJsonFallbackDeserialization() + { + var evt = new EventEnvelope { StreamId = "s-raw", Version = 1, TenantId = "t1", GlobalSequence = 77, Tags = new HashSet { "Patient" }, Data = new TestEventPayload("o-raw", 5m) }; + var rawJson = Newtonsoft.Json.JsonConvert.SerializeObject(evt); + + var env = new CosmosDocumentEnvelope { Id = "e-raw", PartitionKey = "s-raw", TenantId = "t1", Data = rawJson }; + + var iterator = Substitute.For>>(); + var page = Substitute.For>>(); + page.GetEnumerator().Returns(new List> { env }.GetEnumerator()); + page.RequestCharge.Returns(2.0); + + iterator.HasMoreResults.Returns(true, false); + iterator.ReadNextAsync(Arg.Any()).Returns(Task.FromResult(page)); + + _mockEventContainer.GetItemQueryIterator>( + Arg.Any(), + requestOptions: Arg.Any()) + .Returns(iterator); + + var events = await _provider.FetchEventsByTagAsync("Patient", tenantId: "t1", ct: TestContext.Current.CancellationToken); + + events.Count.ShouldBe(1); + events[0].GlobalSequence.ShouldBe(77); + events[0].Tags.ShouldContain("Patient"); + } + // ========================================== // 5. GetStreamHeaderAsync Tests // ==========================================