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
6 changes: 6 additions & 0 deletions test/Grains/TestGrainInterfaces/IGenericInterfaces.cs
Original file line number Diff line number Diff line change
Expand Up @@ -189,6 +189,11 @@ public interface IGenericPingSelf<T> : IGrainWithGuidKey
Task ScheduleDelayedPingToSelfAndDeactivate(IGenericPingSelf<T> target, T t, TimeSpan delay);
}

public interface ILongRunningTaskObserver : IGrainObserver
{
void OnCallStarted(Guid callId);
}

public interface ILongRunningTaskGrain<T> : IGrainWithGuidKey
{
Task<string> GetRuntimeInstanceId();
Expand All @@ -201,6 +206,7 @@ public interface ILongRunningTaskGrain<T> : IGrainWithGuidKey
[AlwaysInterleave]
Task LongWaitGrainCancellationInterleaving(GrainCancellationToken tc, TimeSpan delay, Guid callId);
Task LongWait(CancellationToken tc, TimeSpan delay, Guid callId);
Task LongWaitWithStartNotification(TimeSpan delay, Guid callId, ILongRunningTaskObserver observer, CancellationToken cancellationToken);
[AlwaysInterleave]
Task LongWaitInterleaving(CancellationToken tc, TimeSpan delay, Guid callId);
Task CallOtherLongRunningTask(ILongRunningTaskGrain<T> target, CancellationToken tc, TimeSpan delay, Guid callId);
Expand Down
7 changes: 7 additions & 0 deletions test/Grains/TestGrains/GenericGrains.cs
Original file line number Diff line number Diff line change
Expand Up @@ -785,6 +785,13 @@ public async Task LongWaitGrainCancellation(GrainCancellationToken ct, TimeSpan
}

public Task LongWaitInterleaving(CancellationToken ct, TimeSpan delay, Guid callId) => LongWait(ct, delay, callId);

public Task LongWaitWithStartNotification(TimeSpan delay, Guid callId, ILongRunningTaskObserver observer, CancellationToken cancellationToken)
{
observer.OnCallStarted(callId);
return LongWait(cancellationToken, delay, callId);
}

public async Task LongWait(CancellationToken ct, TimeSpan delay, Guid callId)
{
try
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -67,14 +67,29 @@ protected override void ConfigureTestCluster(InProcessTestClusterBuilder builder
public async Task GrainTaskCancellation(int delay)
{
var grain = fixture.GrainFactory.GetGrain<ILongRunningTaskGrain<bool>>(Guid.NewGuid());
using var cts = new CancellationTokenSource();
var callId = Guid.NewGuid();
var grainTask = grain.LongWait(cts.Token, TimeSpan.FromSeconds(10), callId);
cts.CancelAfter(delay);
await Assert.ThrowsAnyAsync<OperationCanceledException>(() => grainTask);
if (delay > 0)
var observer = new LongRunningTaskObserver();
var observerReference = fixture.GrainFactory.CreateObjectReference<ILongRunningTaskObserver>(observer);
try
{
await WaitForCallCancellation(grain, callId);
using var cts = new CancellationTokenSource();
var callId = Guid.NewGuid();
var grainTask = grain.LongWaitWithStartNotification(TimeSpan.FromSeconds(10), callId, observerReference, cts.Token);
if (delay > 0)
{
// A timer does not guarantee that a new activation has begun executing the request.
await observer.WaitForCallToStart(callId);
}

cts.CancelAfter(delay);
await Assert.ThrowsAnyAsync<OperationCanceledException>(() => grainTask);
if (delay > 0)
{
await WaitForCallCancellation(grain, callId);
}
}
finally
{
fixture.GrainFactory.DeleteObjectReference<ILongRunningTaskObserver>(observerReference);
}
}

Expand Down Expand Up @@ -353,6 +368,19 @@ private async Task WaitForCallCancellation<T>(ILongRunningTaskGrain<T> grain, Gu
Assert.Fail("Did not encounter the expected call id");
}

private sealed class LongRunningTaskObserver : ILongRunningTaskObserver
{
private readonly TaskCompletionSource<Guid> _callStarted = new(TaskCreationOptions.RunContinuationsAsynchronously);

public void OnCallStarted(Guid callId) => _callStarted.TrySetResult(callId);

public async Task WaitForCallToStart(Guid expectedCallId)
{
var callId = await _callStarted.Task.WaitAsync(TimeSpan.FromSeconds(30));
Assert.Equal(expectedCallId, callId);
}
}

/// <summary>
/// Tests that a running interleaving grain operation can be cancelled via CancellationToken.
/// Interleaving requests run concurrently without queueing and should also be cancellable.
Expand Down
Loading