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
Original file line number Diff line number Diff line change
@@ -0,0 +1,244 @@
using System.Collections.Concurrent;
using IntegrationTests;
using JasperFx.Core;
using Microsoft.Extensions.Hosting;
using Shouldly;
using Wolverine;
using Wolverine.Persistence.Durability;
using Wolverine.Postgresql;
using Wolverine.Postgresql.Transport;
using Wolverine.Runtime;
using Wolverine.Tracking;
using Wolverine.Transports;
using Wolverine.Util;

namespace PostgresqlTests.Transport;

/// <summary>
/// Reproduction for https://github.com/JasperFx/wolverine/issues/4822.
///
/// <para>
/// A scheduled message to a global partition parks in the inbox at the EXTERNAL slot's address (GH-4673), so
/// slot ownership is settled when it comes due. A node that never owned the slot has no listening agent for
/// that address and forwards the envelope to the slot. A node that USED to own it still has one -- stopped
/// when the leader moved the slot away, but still registered -- so the scheduled poller handed the envelopes
/// to a listener with no receiver, which threw after the rows had already been committed as Incoming and
/// owned by this live node. Nothing ever recovered them, and every other destination in the same batch was
/// stranded with them.
/// </para>
///
/// <para>
/// A Solo host owns every slot, so stopping one slot's exclusive listener produces a genuine ex-owner without
/// standing up a second node -- the same technique as GH-4700 and GH-4776. The rows are written straight into
/// the inbox and promoted by the host's own scheduled poller, which is the path the report goes through.
/// </para>
/// </summary>
[Collection("Postgresql")]
public class Bug_4822_scheduled_promotion_after_slot_handoff : IAsyncLifetime
{
private IHost _host = null!;
private WolverineRuntime _runtime = null!;

public async ValueTask InitializeAsync()
{
HandoffStepHandler.Received.Clear();

_host = await Host.CreateDefaultBuilder()
.UseWolverine(opts =>
{
opts.Durability.Mode = DurabilityMode.Solo;
opts.Durability.ScheduledJobFirstExecution = 100.Milliseconds();
opts.Durability.ScheduledJobPollingTime = 250.Milliseconds();

opts.UsePostgresqlPersistenceAndTransport(Servers.PostgresConnectionString, "slot4822",
transportSchema: "slot4822_queues")
.AutoProvision()
.AutoPurgeOnStartup();

opts.Discovery.DisableConventionalDiscovery().IncludeType(typeof(HandoffStepHandler));

opts.MessagePartitioning.ByMessage<HandoffStep>(x => x.GroupId.ToString());

opts.MessagePartitioning.GlobalPartitioned(topology =>
{
topology.UseShardedPostgresqlQueues("handoff", 2);
topology.Message<HandoffStep>();
});

// An ordinary exclusive listener with no partitioning behind it. Stopped, it still has no
// receiver to take promoted envelopes, which makes it a dependable way to fail one group
// of a scheduled batch on purpose.
opts.ListenToPostgresqlQueue("handoffplain").ListenWithStrictOrdering();
}).StartAsync();

_runtime = _host.GetRuntime();
}

public async ValueTask DisposeAsync()
{
await _host.StopAsync();
_host.Dispose();
}

private PostgresqlQueue[] theSlots()
{
var transport = _runtime.Options.Transports.GetOrCreate<PostgresqlTransport>();
return [transport.Queues["handoff1"], transport.Queues["handoff2"]];
}

private PostgresqlQueue thePlainQueue()
{
return _runtime.Options.Transports.GetOrCreate<PostgresqlTransport>().Queues["handoffplain"];
}

/// <summary>
/// Shaped like the row a cascaded <c>ScheduledAt</c> leaves behind: Scheduled, owned by any node, parked
/// at its eventual destination and already due. <paramref name="dueAgo"/> orders the rows inside the
/// poller's batch, which is sorted by execution time.
/// </summary>
private async Task<(Guid GroupId, Guid EnvelopeId)> parkScheduledAsync(Uri destination, TimeSpan dueAgo)
{
var message = new HandoffStep(Guid.NewGuid());
var serializer = _runtime.Options.DefaultSerializer;

var envelope = new Envelope(message)
{
Destination = destination,
MessageType = typeof(HandoffStep).ToMessageTypeName(),
Serializer = serializer,
ContentType = serializer.ContentType,
Status = EnvelopeStatus.Scheduled,
OwnerId = TransportConstants.AnyNode,
ScheduledTime = DateTimeOffset.UtcNow.Subtract(dueAgo)
};

envelope.Data = serializer.Write(envelope);

await _runtime.Storage.Inbox.StoreIncomingAsync(envelope);

return (message.GroupId, envelope.Id);
}

private async Task<Envelope?> inboxRowAsync(Guid id)
{
var all = await _runtime.Storage.Admin.AllIncomingAsync();
return all.FirstOrDefault(x => x.Id == id);
}

private string describeInbox(IEnumerable<Envelope> rows)
{
return rows.Select(x => $"{x.Id}@{x.Destination} {x.Status} owner={x.OwnerId}").Join("; ");
}

/// <summary>
/// Against the URIs the topology itself stamps, for the reason GH-4700's twin gives: a lookup that only
/// agrees with hand-built endpoints would let the fix silently do nothing.
/// </summary>
[Fact]
public void only_the_external_slots_themselves_are_global_partition_slots()
{
foreach (var slot in theSlots())
{
_runtime.Endpoints.IsGlobalPartitionSlot(slot.Uri).ShouldBeTrue();

// The companion queue is GlobalPartitionSlotFor's question, not this one
_runtime.Endpoints.IsGlobalPartitionSlot(slot.GlobalPartitionLocalQueueUri!).ShouldBeFalse();
}

_runtime.Endpoints.IsGlobalPartitionSlot(thePlainQueue().Uri).ShouldBeFalse();
_runtime.Endpoints.IsGlobalPartitionSlot(new Uri("postgresql://nowhere/")).ShouldBeFalse();
}

/// <summary>
/// The reported defect. The slot's row comes first in the batch so that, before the fix, the throw it
/// caused also stranded the other slot's row -- the "7 rows for repro1" half of the report.
/// </summary>
[Fact]
public async Task an_ex_owner_forwards_a_due_message_to_the_slot_and_the_rest_of_the_batch_still_runs()
{
var (lost, kept) = (theSlots()[0], theSlots()[1]);

await _runtime.Endpoints.StopListenerAsync(lost, TestContext.Current.CancellationToken);

_runtime.Endpoints.FindListeningAgent(lost.Uri).ShouldNotBeNull(
"Precondition: an ex-owner still has the stopped listening agent registered");

var forLostSlot = await parkScheduledAsync(lost.Uri, 2.Seconds());
var forKeptSlot = await parkScheduledAsync(kept.Uri, 1.Seconds());

// The slot this node kept is unaffected by the one it gave up
await waitForAsync(() => HandoffStepHandler.Received.Any(x => x.Id == forKeptSlot.GroupId),
async () => describeInbox(await _runtime.Storage.Admin.AllIncomingAsync()));

// The given-up slot's message left the inbox here without running here ...
await waitForAsync(async () => await inboxRowAsync(forLostSlot.EnvelopeId) == null,
async () => describeInbox(await _runtime.Storage.Admin.AllIncomingAsync()));

HandoffStepHandler.Received.ShouldNotContain(x => x.Id == forLostSlot.GroupId);

// ... and was handed to the slot rather than dropped: the slot's owner -- this node again, once it
// gets the listener back -- runs it, exactly once.
await _runtime.Endpoints.StartListenerAsync(lost, TestContext.Current.CancellationToken);

await waitForAsync(() => HandoffStepHandler.Received.Any(x => x.Id == forLostSlot.GroupId),
async () => describeInbox(await _runtime.Storage.Admin.AllIncomingAsync()));

await Task.Delay(1.Seconds(), TestContext.Current.CancellationToken);
HandoffStepHandler.Received.Count(x => x.Id == forLostSlot.GroupId).ShouldBe(1);
HandoffStepHandler.Received.Single(x => x.Id == forLostSlot.GroupId).Destination.ShouldBe(lost.Uri);
}

/// <summary>
/// The general hardening. Whatever makes one destination of a promoted batch fail, the rest of the batch
/// still runs, and the failed rows are handed back to any node instead of staying owned by this live one
/// where no recovery would ever look at them.
/// </summary>
[Fact]
public async Task a_destination_that_cannot_take_promoted_envelopes_does_not_strand_its_rows_or_the_batch()
{
var plain = thePlainQueue();
var kept = theSlots()[1];

await _runtime.Endpoints.StopListenerAsync(plain, TestContext.Current.CancellationToken);

var failing = await parkScheduledAsync(plain.Uri, 2.Seconds());
var healthy = await parkScheduledAsync(kept.Uri, 1.Seconds());

await waitForAsync(() => HandoffStepHandler.Received.Any(x => x.Id == healthy.GroupId),
async () => describeInbox(await _runtime.Storage.Admin.AllIncomingAsync()));

// Promoted, then released rather than stranded under this node's id
await waitForAsync(async () => await inboxRowAsync(failing.EnvelopeId) is { Status: EnvelopeStatus.Incoming } row
&& row.OwnerId == TransportConstants.AnyNode,
async () => describeInbox(await _runtime.Storage.Admin.AllIncomingAsync()));

HandoffStepHandler.Received.ShouldNotContain(x => x.Id == failing.GroupId);
}

private static Task waitForAsync(Func<bool> condition, Func<Task<string>> diagnostic)
{
return waitForAsync(() => Task.FromResult(condition()), diagnostic);
}

private static async Task waitForAsync(Func<Task<bool>> condition, Func<Task<string>> diagnostic)
{
var deadline = DateTimeOffset.UtcNow.Add(30.Seconds());
while (DateTimeOffset.UtcNow < deadline)
{
if (await condition()) return;
await Task.Delay(100.Milliseconds());
}

throw new TimeoutException($"Condition never met. Inbox: {await diagnostic()}");
}
}

public record HandoffStep(Guid GroupId);

public static class HandoffStepHandler
{
public static readonly ConcurrentBag<(Guid Id, Uri? Destination)> Received = new();

public static void Handle(HandoffStep message, Envelope envelope) =>
Received.Add((message.GroupId, envelope.Destination));
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,101 @@
using JasperFx.Core;
using Microsoft.Extensions.Hosting;
using NSubstitute;
using Wolverine.Tracking;
using Wolverine;
using Wolverine.Configuration;
using Wolverine.Persistence.Durability;
using Wolverine.Runtime;
using Wolverine.Transports;
using Wolverine.Transports.Sending;
using Xunit;

namespace CoreTests.Bugs;

/// <summary>
/// GH-4822. When handing a group of promoted scheduled envelopes onward fails, <c>EnqueueDirectlyAsync</c>
/// releases the rows back to any node so that recovery retries them. Two things have to be right about WHICH
/// rows and WHERE:
///
/// <para>
/// Only the envelopes that were not handed over yet. One that already left for the slot has had its inbox row
/// retired; releasing it as well would let recovery run it a second time.
/// </para>
///
/// <para>
/// At the address the row was parked under. The slot forward re-addresses each envelope to the slot before
/// sending (GH-4700), but the release matches on id AND <c>received_at</c>, so naming the live destination
/// would match nothing and leave the row owned by this live node -- the very stranding the release exists to
/// prevent.
/// </para>
///
/// <para>
/// The failing send is a substituted sending agent registered for the slot; the stub transport itself never
/// fails a send.
/// </para>
/// </summary>
public class Bug_4822_failed_promotion_releases_the_parked_rows : IAsyncLifetime
{
private static readonly Uri TheSlot = "stub://partition-slot-4822".ToUri();
private static readonly Uri TheCompanionQueue = "local://global-partition-slot-4822".ToUri();

private IHost _host = null!;
private ISendingAgent _slotSender = null!;
private IMessageStore _store = null!;

public async ValueTask InitializeAsync()
{
_host = await Host.CreateDefaultBuilder()
.UseWolverine(opts =>
{
opts.Discovery.DisableConventionalDiscovery().IncludeType<PartitionedRetryMessageHandler>();
opts.PublishMessage<PartitionedRetryMessage>().To(TheSlot);
})
.StartAsync(TestContext.Current.CancellationToken);

var runtime = _host.GetRuntime();
runtime.Endpoints.EndpointFor(TheSlot)!.GlobalPartitionLocalQueueUri = TheCompanionQueue;

// This node does not own the slot (nothing listens to it here), so a promoted envelope parked at the
// companion queue is forwarded to the slot -- through this sender, whose second send fails.
_slotSender = Substitute.For<ISendingAgent>();
_slotSender.Destination.Returns(TheSlot);
_slotSender.EnqueueOutgoingAsync(Arg.Any<Envelope>())
.Returns(ValueTask.CompletedTask, ValueTask.FromException(new DivideByZeroException()));
((EndpointCollection)runtime.Endpoints).StoreSendingAgent(_slotSender);

_store = Substitute.For<IMessageStore>();
_store.Inbox.Returns(Substitute.For<IMessageInbox>());
}

public async ValueTask DisposeAsync()
{
await _host.StopAsync();
_host.Dispose();
}

private Envelope promoted(string name)
{
return new Envelope(new PartitionedRetryMessage(name))
{
Destination = TheCompanionQueue,
Store = _store
};
}

[Fact]
public async Task only_the_envelopes_not_yet_forwarded_are_released_and_at_their_parked_address()
{
var forwarded = promoted("first");
var failed = promoted("second");
var neverTried = promoted("third");

await _host.GetRuntime().EnqueueDirectlyAsync([forwarded, failed, neverTried]);

await _store.Received(1).ReassignIncomingAsync(TransportConstants.AnyNode,
Arg.Is<IReadOnlyList<Envelope>>(released =>
released.Select(x => x.Id).OrderBy(x => x)
.SequenceEqual(new[] { failed.Id, neverTried.Id }.OrderBy(x => x))
&& released.All(x => x.Destination == TheCompanionQueue)));
}
}
33 changes: 33 additions & 0 deletions src/Wolverine/Configuration/EndpointCollection.cs
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,13 @@ public interface IEndpointCollection : IAsyncDisposable
/// </summary>
Uri? GlobalPartitionSlotFor(Uri localQueueUri);

/// <summary>
/// True when this address is itself a global partition's external slot endpoint -- the address a
/// scheduled message to a global partition parks at (GH-4673). See the implementation for why that is
/// not the same question as <see cref="GlobalPartitionSlotFor"/>.
/// </summary>
bool IsGlobalPartitionSlot(Uri address);

ISendingAgent AgentForLocalQueue(string queueName);
Endpoint? EndpointByName(string endpointName);
IListeningAgent? FindListeningAgent(Uri uri);
Expand Down Expand Up @@ -494,6 +501,32 @@ public bool IsSingleNodeListener(Uri address)
return slotUri;
}

private ImHashMap<Uri, bool> _isGlobalPartitionSlot = ImHashMap<Uri, bool>.Empty;

/// <summary>
/// GH-4822. A scheduled message to a global partition parks in the inbox at the external slot's own
/// address (GH-4673), and <see cref="GlobalPartitionSlotFor"/> deliberately answers null for that address
/// -- it is not a companion queue. Ownership still has to be settled for it at promotion time, though: a
/// node that once owned the slot keeps its stopped listening agent registered, so
/// <see cref="FindListenerCircuit"/> hands the envelopes to a listener with no receiver. Cached the same way
/// as the other lookups here, because the scheduled poller asks once per distinct destination per pass.
/// </summary>
public bool IsGlobalPartitionSlot(Uri address)
{
if (_isGlobalPartitionSlot.TryFind(address, out var isSlot))
{
return isSlot;
}

isSlot = address.Scheme != TransportConstants.Local
&& _options.Transports.AllEndpoints()
.Any(x => x.Uri == address && x.GlobalPartitionLocalQueueUri != null);

_isGlobalPartitionSlot = _isGlobalPartitionSlot.AddOrUpdate(address, isSlot);

return isSlot;
}

public IListenerCircuit? FindListenerCircuit(Uri address)
{
if (address.Scheme == TransportConstants.Local)
Expand Down
Loading
Loading