From cbcc13525917f7a00510f6a2f8efbe5e0b05581c Mon Sep 17 00:00:00 2001 From: Christian Date: Mon, 5 Oct 2026 17:27:23 +0200 Subject: [PATCH 1/2] Forward due scheduled messages from a global partition's ex-owner (GH-4822) A scheduled message to a global partition parks in the inbox at the external slot's address. A node that previously owned the slot keeps its stopped listening agent registered, so the scheduled poller handed the promoted envelopes to a listener without a receiver. That threw after the rows were committed as Incoming and owned by a live node, stranding them and every other destination in the same batch. EnqueueDirectlyAsync now settles slot ownership for the slot address itself and forwards to the slot when this node does not own it. Each destination group is also isolated: a failing group is logged and its envelopes are released back to any node so recovery retries them. --- ..._scheduled_promotion_after_slot_handoff.cs | 244 ++++++++++++++++++ .../Configuration/EndpointCollection.cs | 33 +++ .../WolverineRuntime.EnqueueDirectly.cs | 126 ++++++--- 3 files changed, 364 insertions(+), 39 deletions(-) create mode 100644 src/Persistence/PostgresqlTests/Transport/Bug_4822_scheduled_promotion_after_slot_handoff.cs diff --git a/src/Persistence/PostgresqlTests/Transport/Bug_4822_scheduled_promotion_after_slot_handoff.cs b/src/Persistence/PostgresqlTests/Transport/Bug_4822_scheduled_promotion_after_slot_handoff.cs new file mode 100644 index 0000000000..321ba2b74c --- /dev/null +++ b/src/Persistence/PostgresqlTests/Transport/Bug_4822_scheduled_promotion_after_slot_handoff.cs @@ -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; + +/// +/// Reproduction for https://github.com/JasperFx/wolverine/issues/4822. +/// +/// +/// 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. +/// +/// +/// +/// 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. +/// +/// +[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(x => x.GroupId.ToString()); + + opts.MessagePartitioning.GlobalPartitioned(topology => + { + topology.UseShardedPostgresqlQueues("handoff", 2); + topology.Message(); + }); + + // 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(); + return [transport.Queues["handoff1"], transport.Queues["handoff2"]]; + } + + private PostgresqlQueue thePlainQueue() + { + return _runtime.Options.Transports.GetOrCreate().Queues["handoffplain"]; + } + + /// + /// Shaped like the row a cascaded ScheduledAt leaves behind: Scheduled, owned by any node, parked + /// at its eventual destination and already due. orders the rows inside the + /// poller's batch, which is sorted by execution time. + /// + 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 inboxRowAsync(Guid id) + { + var all = await _runtime.Storage.Admin.AllIncomingAsync(); + return all.FirstOrDefault(x => x.Id == id); + } + + private string describeInbox(IEnumerable rows) + { + return rows.Select(x => $"{x.Id}@{x.Destination} {x.Status} owner={x.OwnerId}").Join("; "); + } + + /// + /// 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. + /// + [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(); + } + + /// + /// 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. + /// + [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); + } + + /// + /// 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. + /// + [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 condition, Func> diagnostic) + { + return waitForAsync(() => Task.FromResult(condition()), diagnostic); + } + + private static async Task waitForAsync(Func> condition, Func> 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)); +} diff --git a/src/Wolverine/Configuration/EndpointCollection.cs b/src/Wolverine/Configuration/EndpointCollection.cs index 918538faf0..8e57e7890c 100644 --- a/src/Wolverine/Configuration/EndpointCollection.cs +++ b/src/Wolverine/Configuration/EndpointCollection.cs @@ -26,6 +26,13 @@ public interface IEndpointCollection : IAsyncDisposable /// Uri? GlobalPartitionSlotFor(Uri localQueueUri); + /// + /// 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 . + /// + bool IsGlobalPartitionSlot(Uri address); + ISendingAgent AgentForLocalQueue(string queueName); Endpoint? EndpointByName(string endpointName); IListeningAgent? FindListeningAgent(Uri uri); @@ -494,6 +501,32 @@ public bool IsSingleNodeListener(Uri address) return slotUri; } + private ImHashMap _isGlobalPartitionSlot = ImHashMap.Empty; + + /// + /// GH-4822. A scheduled message to a global partition parks in the inbox at the external slot's own + /// address (GH-4673), and 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 + /// 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. + /// + 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) diff --git a/src/Wolverine/Runtime/WolverineRuntime.EnqueueDirectly.cs b/src/Wolverine/Runtime/WolverineRuntime.EnqueueDirectly.cs index a3e66ba62f..02ad3b33e4 100644 --- a/src/Wolverine/Runtime/WolverineRuntime.EnqueueDirectly.cs +++ b/src/Wolverine/Runtime/WolverineRuntime.EnqueueDirectly.cs @@ -11,51 +11,99 @@ public async ValueTask EnqueueDirectlyAsync(IReadOnlyList envelopes) var groups = envelopes.GroupBy(x => x.Destination ?? TransportConstants.LocalUri).ToArray(); foreach (var group in groups) { - // GH-4700. Has to come before FindListenerCircuit, which answers yes for ANY local:// address. - // A global partition's companion local queue exists on every node by design (GH-3856), so - // "I can build a circuit here" is not the same question as "this slot is mine to run". The - // scheduled poller is deliberately unfiltered -- it takes every due row behind one per-database - // advisory lock, and GH-4645 depends on it promoting rows for destinations it does not serve -- - // so this is where ownership has to be settled. - var slotUri = Endpoints.GlobalPartitionSlotFor(group.Key); - if (slotUri != null && !thisNodeOwnsPartitionSlot(slotUri)) + // GH-4822. Every envelope here was promoted by a scheduled poller that has already committed its + // rows as Incoming and owned by this node, and inbox recovery only ever releases rows owned by a + // node it has proven dead. Letting one destination's failure escape would leave its rows -- and + // those of every destination after it in the batch -- owned by a live node forever, with nothing + // logged against them and nothing dead lettered. Same remedy as GH-3680 on the recovery side: + // hand them back to any node so recovery tries again. + try + { + await enqueueDirectlyAsync(group); + } + catch (Exception e) + { + Logger.LogError(e, + "Error trying to enqueue {Count} promoted scheduled envelopes for {Destination}. Releasing them back to any node so that they are recovered later", + group.Count(), group.Key); + + await releaseToAnyNodeAsync(group); + } + } + } + + private async Task enqueueDirectlyAsync(IGrouping group) + { + // GH-4700. Has to come before FindListenerCircuit, which answers yes for ANY local:// address. + // A global partition's companion local queue exists on every node by design (GH-3856), so + // "I can build a circuit here" is not the same question as "this slot is mine to run". The + // scheduled poller is deliberately unfiltered -- it takes every due row behind one per-database + // advisory lock, and GH-4645 depends on it promoting rows for destinations it does not serve -- + // so this is where ownership has to be settled. + // + // GH-4822. A scheduled message to a global partition parks at the external slot's own address + // (GH-4673), which is not a companion queue, so it needs the same question asked of the slot + // itself. A node that never owned the slot would fall through to the sender branch below and + // forward correctly anyway, but one that USED to own it still has the stopped listening agent + // registered, and FindListenerCircuit hands the envelopes to a listener with no receiver. + var slotUri = Endpoints.GlobalPartitionSlotFor(group.Key) + ?? (Endpoints.IsGlobalPartitionSlot(group.Key) ? group.Key : null); + if (slotUri != null && !thisNodeOwnsPartitionSlot(slotUri)) + { + await forwardToPartitionSlotAsync(group, slotUri); + return; + } + + var listener = Endpoints.FindListenerCircuit(group.Key); + if (listener != null) + { + await listener.EnqueueDirectlyAsync(group); + } + else + { + // For send-only endpoints (e.g. Azure Service Bus topics), + // there is no listener circuit. Send through the sending agent instead. + ISendingAgent sender; + try { - await forwardToPartitionSlotAsync(group, slotUri); - continue; + sender = Endpoints.GetOrBuildSendingAgent(group.Key); + } + catch (UnknownTransportException e) + { + // The envelopes here have already been read out of persistence and + // reassigned to this node, so throwing would both lose the rest of + // this batch and leave the offending rows stranded -- and the poller + // would rediscover them and throw again on every subsequent run. A + // destination whose transport this node cannot resolve is never going + // to become sendable here, so dead letter the envelopes instead. + // See https://github.com/JasperFx/wolverine/issues/3413. + await deadLetterUnknownDestinationAsync(group, e); + return; } - var listener = Endpoints.FindListenerCircuit(group.Key); - if (listener != null) + foreach (var envelope in group) { - await listener.EnqueueDirectlyAsync(group); + await sender.EnqueueOutgoingAsync(envelope); + await retireForwardedInboxRowAsync(envelope); } - else + } + } + + private async Task releaseToAnyNodeAsync(IEnumerable envelopes) + { + foreach (var byStore in envelopes.GroupBy(x => x.Store ?? Storage)) + { + try + { + await byStore.Key.ReassignIncomingAsync(TransportConstants.AnyNode, byStore.ToArray()); + } + catch (Exception e) { - // For send-only endpoints (e.g. Azure Service Bus topics), - // there is no listener circuit. Send through the sending agent instead. - ISendingAgent sender; - try - { - sender = Endpoints.GetOrBuildSendingAgent(group.Key); - } - catch (UnknownTransportException e) - { - // The envelopes here have already been read out of persistence and - // reassigned to this node, so throwing would both lose the rest of - // this batch and leave the offending rows stranded -- and the poller - // would rediscover them and throw again on every subsequent run. A - // destination whose transport this node cannot resolve is never going - // to become sendable here, so dead letter the envelopes instead. - // See https://github.com/JasperFx/wolverine/issues/3413. - await deadLetterUnknownDestinationAsync(group, e); - continue; - } - - foreach (var envelope in group) - { - await sender.EnqueueOutgoingAsync(envelope); - await retireForwardedInboxRowAsync(envelope); - } + // Deliberately not rethrowing, for the same reason retireForwardedInboxRowAsync does not: the + // rest of the batch still deserves its turn, and a stranded row is recoverable by hand. + Logger.LogError(e, + "Error trying to release {Count} un-enqueued promoted envelopes back to any node", + byStore.Count()); } } } From e5bb5512696edfdfe15df3fbd5038d7dba8b8f6a Mon Sep 17 00:00:00 2001 From: Christian Date: Mon, 5 Oct 2026 17:36:15 +0200 Subject: [PATCH 2/2] Release only un-forwarded envelopes at their parked address (GH-4822) When a promoted group fails part-way, envelopes that were already forwarded have had their inbox rows retired. Releasing them as well let recovery run them a second time, so only the envelopes not yet handed over are released now. The release also matched on the envelope's live destination, which the partition slot forward rewrites to the slot before sending. For a row parked at the companion local queue that matched nothing and left it owned by this live node. Stand-in envelopes now carry the parked address. --- ...iled_promotion_releases_the_parked_rows.cs | 101 ++++++++++++++++++ .../WolverineRuntime.EnqueueDirectly.cs | 40 +++++-- 2 files changed, 133 insertions(+), 8 deletions(-) create mode 100644 src/Testing/CoreTests/Bugs/Bug_4822_failed_promotion_releases_the_parked_rows.cs diff --git a/src/Testing/CoreTests/Bugs/Bug_4822_failed_promotion_releases_the_parked_rows.cs b/src/Testing/CoreTests/Bugs/Bug_4822_failed_promotion_releases_the_parked_rows.cs new file mode 100644 index 0000000000..89f9f19d38 --- /dev/null +++ b/src/Testing/CoreTests/Bugs/Bug_4822_failed_promotion_releases_the_parked_rows.cs @@ -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; + +/// +/// GH-4822. When handing a group of promoted scheduled envelopes onward fails, EnqueueDirectlyAsync +/// releases the rows back to any node so that recovery retries them. Two things have to be right about WHICH +/// rows and WHERE: +/// +/// +/// 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. +/// +/// +/// +/// 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 received_at, 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. +/// +/// +/// +/// The failing send is a substituted sending agent registered for the slot; the stub transport itself never +/// fails a send. +/// +/// +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(); + opts.PublishMessage().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(); + _slotSender.Destination.Returns(TheSlot); + _slotSender.EnqueueOutgoingAsync(Arg.Any()) + .Returns(ValueTask.CompletedTask, ValueTask.FromException(new DivideByZeroException())); + ((EndpointCollection)runtime.Endpoints).StoreSendingAgent(_slotSender); + + _store = Substitute.For(); + _store.Inbox.Returns(Substitute.For()); + } + + 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>(released => + released.Select(x => x.Id).OrderBy(x => x) + .SequenceEqual(new[] { failed.Id, neverTried.Id }.OrderBy(x => x)) + && released.All(x => x.Destination == TheCompanionQueue))); + } +} diff --git a/src/Wolverine/Runtime/WolverineRuntime.EnqueueDirectly.cs b/src/Wolverine/Runtime/WolverineRuntime.EnqueueDirectly.cs index 02ad3b33e4..526f6d17af 100644 --- a/src/Wolverine/Runtime/WolverineRuntime.EnqueueDirectly.cs +++ b/src/Wolverine/Runtime/WolverineRuntime.EnqueueDirectly.cs @@ -17,22 +17,31 @@ public async ValueTask EnqueueDirectlyAsync(IReadOnlyList envelopes) // those of every destination after it in the batch -- owned by a live node forever, with nothing // logged against them and nothing dead lettered. Same remedy as GH-3680 on the recovery side: // hand them back to any node so recovery tries again. + // + // Only what was not handed over yet goes back: an envelope already forwarded has had its row retired, + // and releasing it too would let recovery run it a second time. The listener branch hands a group + // over as a whole and cannot say how far it got, so a failure there releases all of it -- the same + // trade GH-3680 makes, and the realistic failure (a listener with no receiver) throws before the + // first envelope. + var handedOver = new HashSet(); try { - await enqueueDirectlyAsync(group); + await enqueueDirectlyAsync(group, handedOver); } catch (Exception e) { + var stranded = group.Where(x => !handedOver.Contains(x)).ToArray(); + Logger.LogError(e, "Error trying to enqueue {Count} promoted scheduled envelopes for {Destination}. Releasing them back to any node so that they are recovered later", - group.Count(), group.Key); + stranded.Length, group.Key); - await releaseToAnyNodeAsync(group); + await releaseToAnyNodeAsync(stranded, group.Key); } } } - private async Task enqueueDirectlyAsync(IGrouping group) + private async Task enqueueDirectlyAsync(IGrouping group, ISet handedOver) { // GH-4700. Has to come before FindListenerCircuit, which answers yes for ANY local:// address. // A global partition's companion local queue exists on every node by design (GH-3856), so @@ -50,7 +59,7 @@ private async Task enqueueDirectlyAsync(IGrouping group) ?? (Endpoints.IsGlobalPartitionSlot(group.Key) ? group.Key : null); if (slotUri != null && !thisNodeOwnsPartitionSlot(slotUri)) { - await forwardToPartitionSlotAsync(group, slotUri); + await forwardToPartitionSlotAsync(group, slotUri, handedOver); return; } @@ -84,18 +93,31 @@ private async Task enqueueDirectlyAsync(IGrouping group) foreach (var envelope in group) { await sender.EnqueueOutgoingAsync(envelope); + handedOver.Add(envelope); await retireForwardedInboxRowAsync(envelope); } } } - private async Task releaseToAnyNodeAsync(IEnumerable envelopes) + /// + /// GH-4822. The release matches on id AND received_at, so it has to name the address the rows were + /// parked under. The slot forward has already re-addressed every envelope it touched to the slot, so + /// stand-ins carry the parked address rather than mutating the live envelopes back -- the same reasoning + /// as . + /// + private async Task releaseToAnyNodeAsync(IReadOnlyList envelopes, Uri parkedAt) { + if (envelopes.Count == 0) return; + foreach (var byStore in envelopes.GroupBy(x => x.Store ?? Storage)) { try { - await byStore.Key.ReassignIncomingAsync(TransportConstants.AnyNode, byStore.ToArray()); + var released = byStore + .Select(x => new Envelope { Id = x.Id, Destination = parkedAt, Store = x.Store }) + .ToArray(); + + await byStore.Key.ReassignIncomingAsync(TransportConstants.AnyNode, released); } catch (Exception e) { @@ -138,7 +160,8 @@ private bool thisNodeOwnsPartitionSlot(Uri slotUri) /// strand the row, which is exactly the GH-4645 data loss. OutgoingMessageBatch also assigns /// Destination itself, so the live value cannot be trusted once the envelope is handed over. /// - private async Task forwardToPartitionSlotAsync(IEnumerable group, Uri slotUri) + private async Task forwardToPartitionSlotAsync(IEnumerable group, Uri slotUri, + ISet handedOver) { ISendingAgent sender; try @@ -162,6 +185,7 @@ private async Task forwardToPartitionSlotAsync(IEnumerable group, Uri envelope.Destination = slotUri; await sender.EnqueueOutgoingAsync(envelope); + handedOver.Add(envelope); await retireForwardedInboxRowAsync(envelope, parkedAt); } }