Skip to content
Closed
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
3 changes: 3 additions & 0 deletions src/MeshWeaver.Connection.Orleans/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,9 @@ app.StartPortalApplication();
- Integrates with [Microsoft Orleans](https://learn.microsoft.com/en-us/dotnet/orleans/overview) client system
- Supports Azure and other cloud providers for clustering

## Placement expiry
A delivery whose message outlives Orleans' response timeout while its target grain is still being placed is cancelled with `OperationCanceledException: Message expired before placement could complete for grain ...` (issue #5731; seen once on memex for an ephemeral `_Activity/compile-*` hub). The message reached no activation, so the sender is told a transient failure (`ErrorType.ShuttingDown`) on the first attempt and the delivery is not re-sent: `RoutingGrain.IsPlacementExpired` (in `MeshWeaver.Hosting.Orleans`) is the classifier and `PlacementExpiryClassificationTest` pins Orleans' wording. Whether the stall itself is general placement latency during cluster membership changes is not evidenced; watch log fingerprint `bb128e1afa6c6449` for recurrence.

## See Also
- [Orleans Documentation](https://learn.microsoft.com/en-us/dotnet/orleans/overview) - Learn more about Orleans clients
- [Main MeshWeaver Documentation](../../Readme.md) - More about mesh connectivity options
102 changes: 101 additions & 1 deletion src/MeshWeaver.Hosting.Orleans/RoutingGrain.cs
Original file line number Diff line number Diff line change
Expand Up @@ -1678,7 +1678,7 @@ internal static IObservable<IMessageDelivery> DeliverToGrainObservable(

// Defer keeps grainCall COLD so every RetryWhen re-subscribe re-invokes it (fresh grain
// reference → fresh activation). Never Observable.FromAsync — see AsynchronousCalls.md.
return Observable.Defer(() => grainCall().ToObservable())
return ObserveGrainCall(grainCall)
.RetryWhen(errors => errors
.Select((ex, i) => (Exception: ex, Attempt: i))
.SelectMany(t =>
Expand Down Expand Up @@ -1796,6 +1796,8 @@ internal static ErrorType ClassifyDeliveryException(
Exception ex, Func<bool>? scopeDisposed = null, bool activationErrorRecorded = true) =>
OrleansRoutingService.IsDirectoryUnstable(ex)
|| IsShutdownShaped(ex)
// The message expired while its target grain was still being placed (#5731): no activation ever saw it.
|| IsPlacementExpired(ex)
|| IsDeactivatedActivation(ex, activationErrorRecorded)
|| IsScopeTeardown(ex, scopeDisposed)
? ErrorType.ShuttingDown
Expand Down Expand Up @@ -2133,4 +2135,102 @@ internal static bool IsDepartedSiloRejection(Exception e) =>
/// <para>A literal of <b>Orleans.Runtime</b>.</para>
/// </summary>
internal const string SupersededSiloGenerationMarker = "The target silo is no longer active";

/// <summary>
/// Orleans' own wording from <c>PlacementService.PlacementWorker.ExecutePlacementAsync</c>:
/// <c>"Message expired before placement could complete for grain {id}."</c>, raised as an
/// <see cref="OperationCanceledException"/> when the message's time-to-live (Orleans'
/// <c>ResponseTimeout</c>) ran out while the grain was still being PLACED (issue #5731). Kept
/// as the verbatim prefix without the grain id, so it matches whichever grain the expiry names.
///
/// <para>A literal of <b>Orleans.Runtime</b>. <c>PlacementExpiryClassificationTest</c> pins it
/// against the shipped assembly, so a re-worded Orleans turns that test red instead of turning
/// <see cref="IsPlacementExpired"/> silently inert.</para>
/// </summary>
internal const string PlacementExpiredMarker = "Message expired before placement could complete";

/// <summary>
/// <b>The delivery's message EXPIRED while Orleans was still placing its target grain - issue
/// #5731.</b> Orleans' <c>PlacementService</c> gives up on a message whose time-to-live has run out,
/// and the caller receives an <see cref="OperationCanceledException"/> carrying
/// <see cref="PlacementExpiredMarker"/>. Seen once on memex (2026-09-24 13:43:20Z): two
/// <c>IMessageHubGrain.DeliverMessage</c> calls to an ephemeral <c>_Activity/compile-*</c> hub.
///
/// <para><b>Why it is a transient verdict for the sender.</b> The message died while being
/// ADDRESSED: no activation existed, so it reached no hub and nothing was half-applied. That is a
/// statement about the cluster (placement or the grain directory was slower than the message's
/// time-to-live, which is what a membership change produces), not about the target, and the
/// sender's recovery is bounded the same way <see cref="IsDepartedSiloRejection"/> describes.
/// Reported as <see cref="ErrorType.Failed"/> it tore down consumers that ride out
/// <see cref="ErrorType.ShuttingDown"/>.</para>
///
/// <para><b>Classification only - the delivery is NOT sent again.</b> Each re-attempt can itself
/// wait a full time-to-live inside placement while holding a routing-pool slot, so six retries
/// under a persistent stall would pin that slot for minutes - the amplification #1172 removed for
/// response timeouts. <see cref="IsTransientFailure"/> therefore does not match it, and the sender
/// is told on the first attempt.</para>
///
/// <para><b>Narrow on purpose.</b> A cancellation is also what a hub raises while it shuts down or
/// times out an activation, and those are different statements. Only the one Orleans sentence
/// qualifies, and only on an <see cref="OperationCanceledException"/> (an
/// <see cref="InvalidOperationException"/> quoting the words stays terminal). It is prose, so
/// <c>PlacementExpiryClassificationTest</c> pins it against the shipped Orleans build: if that
/// test fails after an upgrade, repair the marker and never delete the test. The exception reaches
/// this predicate with its text intact only because <see cref="ObserveGrainCall{T}"/> hands it
/// over - see there.</para>
/// </summary>
/// <param name="ex">The exception the delivery attempt faulted with.</param>
/// <returns><c>true</c> when the message expired before its target grain was placed.</returns>
internal static bool IsPlacementExpired(Exception ex) =>
ExceptionChain.Contains(ex, static e =>
e is OperationCanceledException
&& e.Message.Contains(PlacementExpiredMarker, StringComparison.OrdinalIgnoreCase));

/// <summary>
/// The cold grain call as an observable, with a CANCELLED call's original exception kept intact -
/// issue #5731. Orleans completes a call it cannot address (a message that expired in placement)
/// with the carried <see cref="OperationCanceledException"/>, which <c>ValueTask.AsTask</c> turns
/// into a CANCELLED task that still stores it. Rx's task bridge then replaces that with a fresh
/// <see cref="TaskCanceledException"/> whose text is only "A task was canceled." - so every
/// classifier downstream, and the sender's NACK, lost the one sentence that says what happened.
///
/// <para>This hands over the stored exception instead. Anything else (a fault, a result, a
/// cancellation that stored nothing) passes through unchanged, and the retry around it is
/// unaffected: only the exception's identity changes, never whether it is retried.</para>
/// </summary>
/// <typeparam name="T">The grain call's result.</typeparam>
/// <param name="grainCall">The grain call; re-invoked from scratch on every subscription.</param>
/// <returns>A cold observable of the call's result.</returns>
internal static IObservable<T> ObserveGrainCall<T>(Func<Task<T>> grainCall) =>
Observable.Defer(() => grainCall().ToObservable())
.Catch<T, TaskCanceledException>(canceled => Observable.Throw<T>(OriginalCancellation(canceled)));

/// <summary>
/// The exception a cancelled task stored, or <paramref name="canceled"/> itself when it stored
/// none. The task is already complete, so reading it never blocks - it only re-throws what it
/// holds.
/// </summary>
/// <param name="canceled">The exception Rx made from the cancelled task.</param>
/// <returns>The task's own cancellation, or <paramref name="canceled"/>.</returns>
internal static Exception OriginalCancellation(TaskCanceledException canceled)
{
var task = canceled.Task;
if (task is null || !task.IsCanceled)
return canceled;

try
{
task.GetAwaiter().GetResult();
}
catch (OperationCanceledException original)
{
return original;
}
catch (Exception)
{
// Not a cancellation after all: keep the exception Rx made.
}

return canceled;
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,209 @@
using System;
using System.Collections.Generic;
using System.IO;
using System.Reactive.Concurrency;
using System.Text;
using System.Threading.Tasks;
using MeshWeaver.Messaging;
using Microsoft.Extensions.Logging.Abstractions;
using Xunit;
using MeshWeaver.Fixture;

namespace MeshWeaver.Hosting.Orleans.Test;

/// <summary>
/// Issue #5731 - a delivery whose message EXPIRED while Orleans was still placing its target grain.
///
/// <para>Orleans' PlacementService cancels such a message with an OperationCanceledException that
/// the caller receives as a CANCELLED task. Two things were wrong on our side: Rx's task bridge
/// replaced the stored exception with a message-less TaskCanceledException ("A task was canceled."),
/// so nothing downstream could recognise the condition, and the classifier would not have known it
/// anyway, so the sender was told a terminal Failed and consumers that ride out ShuttingDown tore
/// down. The facts here are pure and need no cluster; the exception text is quoted verbatim from the
/// production log.</para>
/// </summary>
public class PlacementExpiryClassificationTest
{
private static readonly Func<int, TimeSpan> NoBackoff = _ => TimeSpan.Zero;

/// <summary>
/// Verbatim from the production log (#5731, 2026-09-24 13:43:20Z, Orleans.Messaging[100071]).
/// </summary>
private const string PlacementExpiredText =
"Message expired before placement could complete for grain "
+ "messagehub/Signature/Desk/_Activity/compile-20260924134235046d3415968dc0b4b9cad1eddb1be5b8bb6.";

/// <summary>
/// The SHAPE Orleans hands the caller: an async completion that throws an OperationCanceledException
/// ends as a CANCELLED task that still stores it - the same state ValueTask.AsTask produces for a
/// carried cancellation.
/// </summary>
private static async Task<IMessageDelivery> ExpiredBeforePlacement()
{
await Task.Yield();
throw new OperationCanceledException(PlacementExpiredText);
}

private static async Task<IMessageDelivery> CancelledForSomeOtherReason()
{
await Task.Yield();
throw new OperationCanceledException("The operation was canceled.");
}

[Fact]
public async Task APlacementExpiry_ReachesTheClassifierWithItsOwnText_AndIsNotResent()
{
var calls = 0;
Exception? caught = null;

try
{
await RoutingGrain.DeliverToGrainObservable(
grainCall: () =>
{
calls++;
return ExpiredBeforePlacement();
},
grainKey: "messagehub/Signature/Desk/_Activity/compile-1",
deliveryId: "placement-expiry-1",
logger: NullLogger.Instance,
backoff: NoBackoff,
scheduler: Scheduler.Immediate)
.Await(TestContext.Current.CancellationToken);
}
catch (Exception ex)
{
caught = ex;
}

Assert.NotNull(caught);
calls.Should().Be(1,
"the message died while being addressed; sending it again would wait out another full "
+ "time-to-live inside placement while holding a routing-pool slot, so the sender is told "
+ "on the first attempt");
RoutingGrain.IsPlacementExpired(caught).Should().BeTrue(
"Rx's task bridge turns a cancelled task into a message-less TaskCanceledException, and the "
+ "grain call must hand over the exception the task stored instead - otherwise no classifier "
+ "downstream can see the Orleans sentence");
}

[Fact]
public async Task APlacementExpiry_IsNackedAsTransient_NotTerminal()
{
var nacks = new List<(string Message, ErrorType Type)>();

await RoutingGrain.DeliverToGrainRoute(
grainCall: ExpiredBeforePlacement,
grainKey: "messagehub/Signature/Desk/_Activity/compile-1",
addressPath: "Signature/Desk/_Activity/compile-1",
deliveryId: "placement-expiry-2",
postFailureToSender: (message, type) => nacks.Add((message, type)),
logger: NullLogger.Instance,
backoff: NoBackoff,
scheduler: Scheduler.Immediate)
.Await(TestContext.Current.CancellationToken);

nacks.Count.Should().Be(1, "the route answers the sender exactly once");
nacks[0].Type.Should().Be(ErrorType.ShuttingDown,
"the message reached no activation - a statement about the cluster, not the target - and the "
+ "consumers with their own recovery machinery ride out ShuttingDown and tear down on Failed");
nacks[0].Message.Contains(RoutingGrain.PlacementExpiredMarker, StringComparison.OrdinalIgnoreCase)
.Should().BeTrue("the sender's NACK now names what happened instead of 'A task was canceled.'");
}

[Fact]
public void APlacementExpiry_IsClassifiedTransient_ButNeverResent()
{
var expired = new OperationCanceledException(PlacementExpiredText);

RoutingGrain.ClassifyDeliveryException(expired).Should().Be(ErrorType.ShuttingDown);
RoutingGrain.ClassifyPodHubDeliveryException(expired).Should().Be(ErrorType.ShuttingDown,
"the pod-hub leg classifies through the same arms");
RoutingGrain.IsResendableDeliveryFailure(expired).Should().BeFalse(
"classification only: a re-send under a persistent placement stall would hold a routing-pool "
+ "slot for a full time-to-live per attempt");
}

[Fact]
public void ItIsFoundBesideAnotherFaultInAnAggregate()
{
var twoTransports = new AggregateException(
new InvalidOperationException("the NACK's other transport also failed"),
new OperationCanceledException(PlacementExpiredText));

RoutingGrain.ClassifyDeliveryException(twoTransports).Should().Be(ErrorType.ShuttingDown,
"which fault sits at index 0 of an aggregate is a race, so the walk reads the whole graph");
}

[Fact]
public async Task AnyOtherCancellation_StaysTerminal()
{
RoutingGrain.ClassifyDeliveryException(new OperationCanceledException("The operation was canceled."))
.Should().Be(ErrorType.Failed,
"a hub raising a cancellation while it shuts down or times out is a different statement "
+ "from an expiry in placement - only Orleans' one sentence qualifies");

var nacks = new List<(string Message, ErrorType Type)>();

await RoutingGrain.DeliverToGrainRoute(
grainCall: CancelledForSomeOtherReason,
grainKey: "messagehub/Signature/Desk/_Activity/compile-1",
addressPath: "Signature/Desk/_Activity/compile-1",
deliveryId: "placement-expiry-3",
postFailureToSender: (message, type) => nacks.Add((message, type)),
logger: NullLogger.Instance,
backoff: NoBackoff,
scheduler: Scheduler.Immediate)
.Await(TestContext.Current.CancellationToken);

nacks.Count.Should().Be(1);
nacks[0].Type.Should().Be(ErrorType.Failed,
"anything this does not recognise stays terminal, so a genuine defect is still reported as one");
}

[Fact]
public void TheMarkerQuotedByAnythingButACancellation_StaysTerminal()
{
RoutingGrain.ClassifyDeliveryException(new InvalidOperationException(PlacementExpiredText))
.Should().Be(ErrorType.Failed,
"the words alone are not the signal: the type guard is what keeps an unrelated fault that "
+ "quotes them from being demoted to a transient verdict");
}

/// <summary>
/// THE ANTI-INERT PIN. If this fact fails after an Orleans upgrade, IsPlacementExpired has gone
/// SILENTLY INERT - repair the marker, never delete the test. The phrase is a literal of the SHIPPED
/// Orleans.Runtime (PlacementService lives there), not a copy of it.
/// </summary>
[Fact]
public void PlacementExpiredMarker_IsStillTheShippedOrleansWording()
{
var runtimeAssembly = Path.Combine(AppContext.BaseDirectory, "Orleans.Runtime.dll");

File.Exists(runtimeAssembly).Should().BeTrue(
$"this fact reads Orleans' own string literals out of {runtimeAssembly}; without the file it "
+ "would pass having verified nothing, which is precisely the failure mode it exists to prevent");

ShippedLiteralPresent(runtimeAssembly, RoutingGrain.PlacementExpiredMarker).Should().BeTrue(
$"'{RoutingGrain.PlacementExpiredMarker}' is the phrase RoutingGrain.IsPlacementExpired matches "
+ "an expiry in Orleans' PlacementService on. If Orleans reworded it, a delivery that expired "
+ "before its target was placed is NACK'd as permanent again (#5731). REPAIR THE MARKER; do not "
+ "delete this");
}

/// <summary>
/// Reads the #US heap at BOTH byte alignments - a managed string literal's UTF-16 payload follows a
/// compressed-integer length, so a blob can begin at an ODD file offset and decoding from offset 0
/// alone would miss a phrase that is genuinely present (a false red on the one assertion whose value
/// is being believed when it fires). OrdinalIgnoreCase, matching the classifier.
/// </summary>
private static bool ShippedLiteralPresent(string assemblyPath, string phrase)
{
var raw = File.ReadAllBytes(assemblyPath);

return Encoding.Unicode.GetString(raw)
.Contains(phrase, StringComparison.OrdinalIgnoreCase)
|| Encoding.Unicode.GetString(raw, 1, raw.Length - 1)
.Contains(phrase, StringComparison.OrdinalIgnoreCase);
}
}
Loading
Loading