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
104 changes: 104 additions & 0 deletions src/Testing/CoreTests/Acceptance/batching_processor_build_race.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,104 @@
using JasperFx.Core;
using Microsoft.Extensions.Hosting;
using Microsoft.Extensions.DependencyInjection;
using Wolverine.Runtime;
using Wolverine.Runtime.Batching;
using Wolverine.Tracking;
using Xunit;

namespace CoreTests.Acceptance;

/// <summary>
/// GH-4167. BatchingOptions.BuildHandler is lazy init on the message-handling path, reached through
/// HandlerPipeline's LightweightCache&lt;Type, IExecutor&gt; whose indexer does not lock — two concurrent
/// misses each invoke the factory and each returns its own instance.
///
/// A duplicate stateless executor is harmless. A duplicate BatchingProcessor is not: each owns a
/// separate BatchingChannel buffer, its own flush Timer, and two Blocks with live worker tasks. Two
/// instances means the batch silently splits, and the loser is never returned to anyone so it is
/// never disposed — its timer and worker tasks leak for the life of the process.
/// </summary>
public class batching_processor_build_race
{
private static Task<IHost> hostAsync(TimeSpan triggerTime, string queueName)
{
return Host.CreateDefaultBuilder()
.UseWolverine(opts =>
{
opts.BatchMessagesOf<ExpiryItem>(batching =>
{
batching.TriggerTime = triggerTime;
batching.LocalExecutionQueueName = queueName;
});
}).StartAsync();
}

/// <summary>
/// The direct expression of the race, and the one that fails without the fix regardless of which
/// JasperFx version is referenced. The end-to-end test below can only fail once Block stops running
/// continuations inline on the publisher (jasperfx#714), because that inline execution serializes
/// the first messages onto one thread and closes the window.
/// </summary>
[Fact]
public async Task concurrent_BuildHandler_yields_one_shared_processor()
{
using var host = await hostAsync(1.Seconds(), "race_direct");

var runtime = (WolverineRuntime)host.Services.GetRequiredService<IWolverineRuntime>();
var options = runtime.Options.BatchDefinitions.Single(x => x.ElementType == typeof(ExpiryItem));

const int racers = 8;
var gate = new Barrier(racers);

var handlers = await Task.WhenAll(Enumerable.Range(0, racers).Select(_ => Task.Run(() =>
{
gate.SignalAndWait();
return options.BuildHandler(runtime);
})));

// Reference equality: every racer must have been handed the SAME processor. Without the fix
// each concurrent miss builds and returns its own, each with its own buffer and flush timer.
handlers.Distinct().Count().ShouldBe(1,
$"BuildHandler produced {handlers.Distinct().Count()} distinct BatchingProcessor instances " +
"across concurrent callers; every extra one owns an orphaned Timer and worker tasks.");
}

/// <summary>
/// End-to-end consequence: the two members of a concurrent first wave must land in one batch.
/// NOTE: this passes against JasperFx 2.56.0 even without the fix, because Block's inline
/// continuations serialize the first two messages. It becomes a real guard once jasperfx#714 ships.
/// </summary>
[Fact]
public async Task concurrent_first_messages_still_assemble_a_single_batch()
{
ExpiryItemHandler.Clear();

using var host = await hostAsync(500.Milliseconds(), "race_items");

var session = await host.TrackActivity()
.Timeout(30.Seconds())
.WaitForMessageToBeReceivedAt<ExpiryItem[]>(host)
.ExecuteAndWaitAsync((Func<IMessageContext, Task>)(async c =>
{
var gate = new Barrier(2);

await Task.WhenAll(
Task.Run(async () =>
{
gate.SignalAndWait();
await c.PublishAsync(new ExpiryItem("one"));
}),
Task.Run(async () =>
{
gate.SignalAndWait();
await c.PublishAsync(new ExpiryItem("two"));
}));
}));

var batches = session.Executed.MessagesOf<ExpiryItem[]>().ToArray();

batches.Length.ShouldBe(1,
"The two members were split across separate BatchingProcessor instances: " +
string.Join(" | ", batches.Select(b => "[" + string.Join(",", b.Select(x => x.Name)) + "]")));
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -75,9 +75,18 @@ public async Task batched_handler_actually_executes()
opts.BatchMessagesOf<ItemDeleted3399>();
}).StartAsync(cancellationToken: TestContext.Current.CancellationToken);

// GH-4167: wait for EXECUTION, not receipt. A WaitFor* condition REPLACES quiescence rather
// than adding to it, so WaitForMessageToBeReceivedAt released this session the moment the batch
// envelope arrived -- before any handler ran -- and the assertion below raced it. That passed
// only because JasperFx ran Block continuations inline on the publisher, which made receive and
// execute effectively atomic; it failed on every run once that was fixed (jasperfx#714).
//
// The count is 2 because that is the whole point of this fixture: under Separated behavior the
// batched array has TWO chains (TelemetryHandler3399 and OtherDeletedHandler3399). Waiting on
// one execution just swaps a receive/execute race for a which-chain-won race.
await host.TrackActivity()
.Timeout(30.Seconds())
.WaitForMessageToBeReceivedAt<ItemDeleted3399[]>(host)
.WaitForExecutionOf<ItemDeleted3399[]>(2)
.SendMessageAndWaitAsync(new ItemDeleted3399(Guid.NewGuid()));

TelemetryHandler3399.Batched.ShouldBeGreaterThan(0);
Expand Down
27 changes: 24 additions & 3 deletions src/Wolverine/Runtime/Batching/BatchingOptions.cs
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,11 @@ namespace Wolverine.Runtime.Batching;

public class BatchingOptions : IAsyncDisposable
{
private IMessageHandler _handler = null!;
private volatile IMessageHandler _handler = null!;

// GH-4167. BuildHandler runs on the message-handling path, not at bootstrap, so it has to be
// safe for concurrent first-touch. See the note there for why a duplicate is not benign.
private readonly object _handlerLock = new();

// CloseAndBuildAs over DefaultMessageBatcher<elementType> and
// ProcessorBuilder<ElementType> closes a generic over the user-supplied
Expand Down Expand Up @@ -179,8 +183,25 @@ internal IMessageHandler BuildHandler(WolverineRuntime runtime)
{
if (_handler != null) return _handler;

var builder = typeof(ProcessorBuilder<>).CloseAndBuildAs<IProcessorBuilder>(ElementType);
_handler = builder.Build(runtime, Batcher, this);
// GH-4167. This is lazy init on the message-handling path, and the caller reaches it through
// HandlerPipeline's LightweightCache<Type, IExecutor>, whose indexer does NOT lock: two
// concurrent misses each invoke the factory AND each returns its own instance. A duplicate
// stateless executor is harmless, which is why that cache is fine everywhere else. A duplicate
// BatchingProcessor is not -- each one owns a separate BatchingChannel buffer, its own flush
// Timer, and two Blocks with live worker tasks. Two instances means the batch silently splits
// across them, and the loser is never handed back to anyone, so it is never disposed: its timer
// and worker tasks leak for the life of the process.
//
// Reproduced deterministically once JasperFx stopped running Block continuations inline on the
// publisher (jasperfx#714) -- that inline execution had been serializing the first messages onto
// one thread and hiding this. The race was always reachable under genuinely concurrent traffic.
lock (_handlerLock)
{
if (_handler != null) return _handler;

var builder = typeof(ProcessorBuilder<>).CloseAndBuildAs<IProcessorBuilder>(ElementType);
_handler = builder.Build(runtime, Batcher, this);
}

return _handler;
}
Expand Down
Loading