Skip to content
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,11 @@ static class ChannelExtensions {
public async Task Read(ProcessElement<T> process, CancellationToken cancellationToken) {
try {
while (!cancellationToken.IsCancellationRequested) {
var element = await channel.Reader.ReadAsync(cancellationToken).NoContext();
// Not cancellable on purpose: a bounded channel parks such a read on an operation it keeps
// and reuses, where a cancellable one allocates per wake — per event, on a caught-up
// subscription. So an idle reader is woken only by the channel completing, never by the
// token, which must therefore not be cancelled before the channel is completed (see Stop).
var element = await channel.Reader.ReadAsync(CancellationToken.None).NoContext();
await process(element, cancellationToken).NoContext();
}
} catch (OperationCanceledException) {
Expand All @@ -39,6 +43,7 @@ public async ValueTask Stop(
Task[] readers,
Func<CancellationToken, ValueTask>? finalize = null
) {
// First, before anything cancels cts: completion is the only thing that wakes an idle Read loop.
channel.Writer.TryComplete();

var incompleteReaders = readers.Where(r => !r.IsCompleted).ToArray();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -143,10 +143,10 @@ async Task CommitInternal(CommitPosition position, bool force, CancellationToken
position.LogContext?.CommittingPosition(position);
await _commitCheckpoint(new(_subscriptionId, position.Position), force, cancellationToken).NoContext();
_lastCommit = position;
_positions.RemoveWhere(x => x.Sequence <= position.Sequence);
_positions.RemoveUpTo(position.Sequence);
} catch (OperationCanceledException) {
await _commitCheckpoint(new(_subscriptionId, position.Position), true, default).NoContext();
_positions.RemoveWhere(x => x.Sequence <= position.Sequence);
_positions.RemoveUpTo(position.Sequence);
} catch (Exception e) {
position.LogContext?.UnableToCommitPosition(position, e);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,22 +19,35 @@ public CommitPosition FirstBeforeGap()
};

CommitPosition Get() {
var result = this
.Zip(this.Skip(1), (position1, position2) => (position1, position2))
.FirstOrDefault(tup => tup.position1.Sequence + 1 != tup.position2.Sequence);
var previous = default(CommitPosition);
bool first = true;

if (result == default) return Max;
// The set is ordered by Sequence; the struct enumerator avoids the boxed IEnumerable<T> one
foreach (var current in this) {
if (!first && previous.Sequence + 1 != current.Sequence) {
SubscriptionsEventSource.Log.CheckpointGapDetected(previous, current);

SubscriptionsEventSource.Log.CheckpointGapDetected(result.position1, result.position2);
return result.position1;
return previous;
}

previous = current;
first = false;
}

return Max;
}

/// <summary>
/// Removes all positions with a sequence up to and including the given one. The set is ordered
/// by sequence, so these are always a prefix; no predicate (and no closure) needed.
/// </summary>
internal void RemoveUpTo(ulong sequence) {
// Min is default on an empty set, so the Count check is what stops the loop
while (Count > 0 && Min.Sequence <= sequence) Remove(Min);
}

class PositionsComparer : IComparer<CommitPosition> {
[MethodImpl(MethodImplOptions.AggressiveInlining)]
public int Compare(CommitPosition x, CommitPosition y) {
if (x.Sequence == y.Sequence) return 0;

return x.Sequence > y.Sequence ? 1 : -1;
}
public int Compare(CommitPosition x, CommitPosition y) => x.Sequence.CompareTo(y.Sequence);
}
}
58 changes: 39 additions & 19 deletions src/Core/src/Eventuous.Subscriptions/EventSubscription.cs
Original file line number Diff line number Diff line change
Expand Up @@ -237,9 +237,42 @@ void ReportConnectionDropped(Session session, Failure failure) {
string GetActivityName(string? messageType)
=> _activityNames.GetOrAdd(messageType ?? "", static (type, prefix) => prefix + type, _activityNamePrefix);

protected ValueTask Handler(IMessageConsumeContext context)
=> context.Message == null
? HandleWithoutPayload(context, context is AsyncConsumeContext ? AcknowledgeAsyncContext : null)
: HandleWithPayload(context);

static readonly Acknowledge AcknowledgeAsyncContext = static context => ((AsyncConsumeContext)context).Acknowledge();

/// <summary>
/// A context without a payload is ignored and acknowledged without entering the pipe, so it skips
/// what only the pipe needs: the logging scope and the activity. Hot, since checkpoint-reached
/// contexts arrive payload-less.
/// </summary>
private protected async ValueTask HandleWithoutPayload(IMessageConsumeContext context, Acknowledge? acknowledge) {
// ReSharper disable once NullCoalescingConditionIsAlwaysNotNullAccordingToAPIContract
Logger.Current ??= Log;

Log.MessageReceived(context);

try {
context.Ignore(SubscriptionId);

if (acknowledge != null) await acknowledge(context).NoContext();
} catch (OperationCanceledException e) when (context.CancellationToken.IsCancellationRequested) {
Log.MessageIgnoredWhenStopping(e);
} catch (Exception e) { context.Nack(SubscriptionId, e); }

if (context.HasFailed() && Options.ThrowOnError) {
var exception = context.HandlingResults.GetException();

throw new SubscriptionException(context.Stream, context.MessageType, context.Message, exception ?? new InvalidOperationException());
}
}

// ReSharper disable once CognitiveComplexity
// ReSharper disable once CyclomaticComplexity
protected async ValueTask Handler(IMessageConsumeContext context) {
async ValueTask HandleWithPayload(IMessageConsumeContext context) {
// Use KeyValuePair array instead of Dictionary for 5x speedup and 3x less allocation
var scope = new KeyValuePair<string, object>[] {
new("SubscriptionId", SubscriptionId),
Expand All @@ -251,10 +284,7 @@ protected async ValueTask Handler(IMessageConsumeContext context) {
Logger.Current ??= Log;

using (Log.Logger.BeginScope(scope)) {
// No activity for payload-less contexts: they are ignored and acknowledged below without
// entering the pipe, so an activity would never be started or disposed on the async path —
// a pure allocation leak, hot since checkpoint-reached contexts arrive payload-less.
var activity = EventuousDiagnostics.Enabled && context.Message != null
var activity = EventuousDiagnostics.Enabled
? SubscriptionActivity.Create(
GetActivityName(context.MessageType),
ActivityKind.Internal,
Expand All @@ -269,23 +299,13 @@ protected async ValueTask Handler(IMessageConsumeContext context) {
Log.MessageReceived(context);

try {
if (context.Message != null) {
if (activity != null) {
context.ParentContext = activity.Context;
if (activity != null) {
context.ParentContext = activity.Context;

if (isAsync) { context.Items.AddItem(ContextItemKeys.Activity, activity); }
}

await Pipe.Send(context).NoContext();
if (isAsync) { context.Items.AddItem(ContextItemKeys.Activity, activity); }
}
else {
context.Ignore(SubscriptionId);

if (isAsync) {
var asyncContext = context as AsyncConsumeContext;
await asyncContext!.Acknowledge().NoContext();
}
}
await Pipe.Send(context).NoContext();

if (context.WasIgnored() && activity != null) activity.ActivityTraceFlags = ActivityTraceFlags.None;
} catch (OperationCanceledException e) when (context.CancellationToken.IsCancellationRequested) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -113,7 +113,17 @@ protected async ValueTask HandleInternal(SubscriptionRun run, IMessageConsumeCon
Logger.Current = Log;

var checkpointedRun = (CheckpointedRun)run;
var ctx = new AsyncConsumeContext(context, checkpointedRun.AckMessage, checkpointedRun.NackMessage);

// Never reaches the pipe, so nothing needs the wrapper: the run's own ack takes the context as is.
if (context.Message == null) {
// ReSharper disable once NullCoalescingConditionIsAlwaysNotNullAccordingToAPIContract
context.LogContext ??= Log;
await HandleWithoutPayload(context, checkpointedRun.AckMessage).NoContext();

return;
}

var ctx = new AsyncConsumeContext(context, checkpointedRun.AckMessage, checkpointedRun.NackMessage);
await Handler(ctx).NoContext();
} catch (OperationCanceledException e) when (context.CancellationToken.IsCancellationRequested) {
context.LogContext.MessageHandlingFailed(Options.SubscriptionId, context, e);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,8 +24,17 @@ public AsyncHandlingFilter(uint concurrencyLimit, uint bufferSize = 10) {
_worker = new(Channel.CreateBounded<WorkerTask>(options), DelayedConsume, (int)concurrencyLimit);
}

// Not async on purpose: an async method's ExecutionContext is restored when it returns, so a logger
// context set inside one never reaches the reader loop, and every message pays for a new context.
// Set here it stays with the loop, and is written again only when the next message brings another.
static ValueTask DelayedConsume(WorkerTask workerTask, CancellationToken ct) {
Logger.Current = workerTask.Context.LogContext;

return Consume(workerTask, ct);
}

// ReSharper disable once CognitiveComplexity
static async ValueTask DelayedConsume(WorkerTask workerTask, CancellationToken ct) {
static async ValueTask Consume(WorkerTask workerTask, CancellationToken ct) {
var ctx = workerTask.Context;

using var activity = ctx.Items.GetItem<Activity>(ContextItemKeys.Activity)?.Start();
Expand All @@ -48,8 +57,6 @@ static async ValueTask DelayedConsume(WorkerTask workerTask, CancellationToken c
ctx.CancellationToken = cts.Token;
}

Logger.Current = ctx.LogContext;

try {
try {
await workerTask.Filter.Value.Send(ctx, workerTask.Filter.Next).NoContext();
Expand Down
24 changes: 19 additions & 5 deletions src/Core/src/Eventuous.Subscriptions/Filters/TracingFilter.cs
Original file line number Diff line number Diff line change
Expand Up @@ -19,23 +19,37 @@ public TracingFilter(string consumerName) {
_defaultTags = [.. tags, .. EventuousDiagnostics.Tags];
}

protected override async ValueTask Send(IMessageConsumeContext context, LinkedListNode<IConsumeFilter>? next) {
if (context.Message == null || next == null) return;
protected override ValueTask Send(IMessageConsumeContext context, LinkedListNode<IConsumeFilter>? next) {
if (context.Message == null || next == null) return default;

// The subscription's own activity is reused, not owned: disposing it would stop it before the
// subscription is done with it, so only an activity started here gets disposed.
var reuseCurrent = Activity.Current?.Context == context.ParentContext;

using var started = reuseCurrent
var created = reuseCurrent
? null
: SubscriptionActivity.Start(
: SubscriptionActivity.Create(
$"{Constants.Components.Consumer}.{context.SubscriptionId}/{context.MessageType}",
ActivityKind.Consumer,
context,
_defaultTags
);

var activity = reuseCurrent ? Activity.Current : started;
var reused = reuseCurrent ? Activity.Current : null;

// Nobody listening, or sampled out: nothing to record after the next filter, so its task is returned
// as is rather than awaited.
if (reused == null && created == null) return next.Value.Send(context, next.Next);

return SendTraced(context, next, reused, created);
}

static async ValueTask SendTraced(IMessageConsumeContext context, LinkedListNode<IConsumeFilter> next, Activity? reused, Activity? created) {
// Started here, not by the caller: this method's execution context is restored when it returns, so the
// started activity doesn't stay current for whoever called the filter.
using var started = created?.Start();

var activity = reused ?? started;

if (activity?.IsAllDataRequested == true && context is AsyncConsumeContext asyncConsumeContext) {
activity.SetContextTags(context)?.SetTag(TelemetryTags.Eventuous.Partition, asyncConsumeContext.PartitionId);
Expand Down
5 changes: 3 additions & 2 deletions src/Core/src/Eventuous.Subscriptions/Handlers/EventHandler.cs
Original file line number Diff line number Diff line change
Expand Up @@ -57,8 +57,9 @@ ValueTask<EventHandlingStatus> NoHandler() {
}
}

public override async ValueTask<EventHandlingStatus> HandleEvent(IMessageConsumeContext context)
=> !_handlersMap.TryGetValue(context.Message!.GetType(), out var handler) ? EventHandlingStatus.Ignored : await handler(context).NoContext();
// Not async: there is nothing to do after the handler, so its task is returned as is rather than awaited
public override ValueTask<EventHandlingStatus> HandleEvent(IMessageConsumeContext context)
=> !_handlersMap.TryGetValue(context.Message!.GetType(), out var handler) ? Ignored : handler(context);

public override string ToString() {
var sb = new StringBuilder();
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,77 @@
// Copyright (C) Eventuous HQ OÜ. All rights reserved
// Licensed under the Apache License, Version 2.0.

using System.Collections.Concurrent;
using Eventuous.Subscriptions;
using Eventuous.Subscriptions.Context;
using Eventuous.Subscriptions.Filters;
using Eventuous.Subscriptions.Logging;
using Shouldly;

namespace Eventuous.Tests.Subscriptions;

/// <summary>
/// The filter's reader keeps the logger context between messages instead of setting it anew for each, so
/// what's pinned here is that a message still never runs under the context of the one before it — the case
/// of a pipe shared by two subscriptions, whose messages interleave on one reader.
/// </summary>
public class AsyncHandlingFilterLogContextTests {
[Test]
public async Task Handler_sees_the_log_context_of_its_own_message() {
var handler = new RecordingHandler();

await using var pipe = new ConsumePipe().AddDefaultConsumer(handler).AddFilterFirst(new AsyncHandlingFilter(1));

var first = Logger.CreateContext("first", null);
var second = Logger.CreateContext("second", null);

// Repeats and switches both: a context kept from the previous message must be neither lost nor stale.
LogContext[] sent = [first, first, second, first, second, second];

for (var i = 0; i < sent.Length; i++) {
await pipe.Send(CreateContext(i, sent[i]));
}

var handled = await Wait.Until(() => handler.Seen.Count == sent.Length, TimeSpan.FromSeconds(5));

handled.ShouldBeTrue($"{handler.Seen.Count} of {sent.Length} messages handled");

// One reader, so the order is the order sent.
var seen = handler.Seen.ToArray();

for (var i = 0; i < sent.Length; i++) {
seen[i].Message.ShouldBeSameAs(sent[i], $"message {i} arrived out of order");
seen[i].Current.ShouldBeSameAs(sent[i], $"message {i} was handled under another message's log context");
}
}

static AsyncConsumeContext CreateContext(int position, LogContext logContext) {
var context = new MessageConsumeContext(
Guid.NewGuid().ToString(),
"TestEvent",
"application/json",
"test-stream",
(ulong)position,
(ulong)position,
(ulong)position,
(ulong)position,
DateTime.UtcNow,
new { Number = position },
new(),
logContext.SubscriptionId,
CancellationToken.None
) { LogContext = logContext };

return new(context, _ => default, (_, _) => default);
}

sealed class RecordingHandler : BaseEventHandler {
public ConcurrentQueue<(LogContext Message, LogContext Current)> Seen { get; } = new();

public override ValueTask<EventHandlingStatus> HandleEvent(IMessageConsumeContext context) {
Seen.Enqueue((context.LogContext, Logger.Current));

return new(EventHandlingStatus.Success);
}
}
}
Loading
Loading