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
41 changes: 37 additions & 4 deletions src/DaemonTests/Aggregations/side_effects_in_aggregations.cs
Original file line number Diff line number Diff line change
Expand Up @@ -464,12 +464,30 @@ public class SideEffects2: IRevisioned

public class RecordingMessageOutbox: IMessageOutbox
{
public readonly List<RecordingMessageBatch> Batches = new();
private readonly List<RecordingMessageBatch> _batches = new();

// The daemon raises side effects for the slices in a batch concurrently, so both of the
// collections in this harness have to be synchronized -- an unguarded List.Add silently drops
// batches and messages, which is what made blue_green_side_effect_gate look flaky (#5065).
// The sibling copy in Composites/composite_rebuild_suppresses_side_effects.cs already locks.
public IReadOnlyList<RecordingMessageBatch> Batches
{
get
{
lock (_batches)
{
return _batches.ToList();
}
}
}

public ValueTask<IMessageBatch> CreateBatch(DocumentSessionBase session)
{
var batch = new RecordingMessageBatch();
Batches.Add(batch);
lock (_batches)
{
_batches.Add(batch);
}

return new ValueTask<IMessageBatch>(batch);
}
Expand All @@ -479,7 +497,18 @@ public record TenantMessage(string tenantId, object message);

public class RecordingMessageBatch: IMessageBatch
{
public readonly List<TenantMessage> Messages = new();
private readonly List<TenantMessage> _messages = new();

public IReadOnlyList<TenantMessage> Messages
{
get
{
lock (_messages)
{
return _messages.ToList();
}
}
}

public Task AfterCommitAsync(IDocumentSession session, IChangeSet commit, CancellationToken token)
{
Expand All @@ -498,7 +527,11 @@ public Task BeforeCommitAsync(IDocumentSession session, IChangeSet commit, Cance
public bool BeforeCommitWasCalled { get; set; }
public ValueTask PublishAsync<T>(T message, string tenantId)
{
Messages.Add(new TenantMessage(tenantId, message));
lock (_messages)
{
_messages.Add(new TenantMessage(tenantId, message));
}

return new ValueTask();
}
}
17 changes: 17 additions & 0 deletions src/Marten/Events/Aggregation/IMessageBatch.cs
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,23 @@

namespace Marten.Events.Aggregation;

/// <summary>
/// A batch of messages published by projection side effects within a single projection update.
/// </summary>
/// <remarks>
/// <para>
/// <b>Implementations must be thread safe.</b> The async daemon raises side effects for the event
/// slices in one batch concurrently, so <see cref="IMessageSink.PublishAsync{T}" /> is called from
/// several threads against the same <see cref="IMessageBatch" /> instance -- measured at up to 8
/// simultaneous callers across 10 threads for a single-stream projection catching up over 20
/// streams. An implementation that appends to an unsynchronized collection silently loses
/// messages under load; see #5065, where exactly that made a daemon test look flaky.
/// </para>
/// <para>
/// The same applies to any collection an <see cref="IMessageOutbox" /> uses to track the batches
/// it hands out, since <see cref="IMessageOutbox.CreateBatch" /> is reached from those threads too.
/// </para>
/// </remarks>
public interface IMessageBatch: IMessageSink, IChangeListener
{

Expand Down
Loading