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
Expand Up @@ -106,19 +106,25 @@ private async Task WaitForTask(uint taskId)
{
while (true)
{
// the completion check and the signal reset must happen in the same
// lock section. `Complete()` adds the task id under this lock before
// it calls `_signal.Set()`, so a completion that is not yet visible
// to the check can only signal after the reset. The wakeup cannot be
// lost.
lock (_sync)
{
if (_completed.Contains(taskId))
{
return;
}

// we are waiting for completion of the current task
// so we force the `TryDispatchOrCompleteUnsafe` to seek completion
// even though the _work backlog might still have unprocessed
// tasks.
TryDispatchOrCompleteUnsafe(isWaitingForTaskCompletion: true);
}

// we are waiting for completion of the current task
// so we force the `TryDispatchOrComplete` to seek completion
// even though the _work backlog might still have unprocessed
// tasks.
TryDispatchOrComplete(isWaitingForTaskCompletion: true);
await TryPauseAsync().ConfigureAwait(false);
}
}
Expand Down Expand Up @@ -187,7 +193,7 @@ private void HandleError(Exception exception)
}
}

private void TryDispatchOrComplete(bool isWaitingForTaskCompletion = false)
private void TryDispatchOrComplete()
{
if (_isCompleted)
{
Expand All @@ -196,33 +202,38 @@ private void TryDispatchOrComplete(bool isWaitingForTaskCompletion = false)

lock (_sync)
{
if (_isCompleted)
{
return;
}
TryDispatchOrCompleteUnsafe(isWaitingForTaskCompletion: false);
}
}

if (!isWaitingForTaskCompletion)
{
isWaitingForTaskCompletion = _work is { HasRunningTasks: true, IsEmpty: true };
}
private void TryDispatchOrCompleteUnsafe(bool isWaitingForTaskCompletion)
{
if (_isCompleted)
{
return;
}

if (!isWaitingForTaskCompletion)
{
isWaitingForTaskCompletion = _work is { HasRunningTasks: true, IsEmpty: true };
}

var hasWork = !_work.IsEmpty || !_serial.IsEmpty;
var hasWork = !_work.IsEmpty || !_serial.IsEmpty;

if (isWaitingForTaskCompletion)
{
_signal.Reset();
if (isWaitingForTaskCompletion)
{
_signal.Reset();

if (Interlocked.CompareExchange(ref _hasBatches, 0, 1) == 1)
{
_batchDispatcher.BeginDispatch(_ct);
}
if (Interlocked.CompareExchange(ref _hasBatches, 0, 1) == 1)
{
_batchDispatcher.BeginDispatch(_ct);
}
else
}
else
{
if (!hasWork && _pendingBatches.Count == 0)
{
if (!hasWork && _pendingBatches.Count == 0)
{
_isCompleted = true;
}
_isCompleted = true;
}
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -69,6 +69,71 @@ public async Task Ensure_Mutations_Child_Fields_Are_Scoped_To_Its_Parent()
""");
}

[Fact]
public async Task Execute_Should_CompleteAllMutations_When_ChildResolverYieldsUnderConcurrency()
{
// arrange
// A mutation root field is executed as a serial task and awaited by the work
// scheduler. Its payload exposes a single async child field that yields, so the
// child resolver completes on a thread-pool thread while the scheduler loop is
// waiting for the serial task. Firing many of these concurrently exercises the
// narrow window where a completion signal can race the scheduler's wait, which
// would leave a request hung forever.
var executor =
await new ServiceCollection()
.AddGraphQL()
.AddQueryType(d => d.Field("noop").Resolve("noop"))
.AddMutationType<RaceMutation>()
.BuildRequestExecutorAsync(cancellationToken: TestContext.Current.CancellationToken);

const string mutation =
"""
mutation {
doWork {
value
}
}
""";

const int rounds = 40;
const int requestsPerRound = 200;

// act
// assert
// Each round must complete within the guard; a lost wakeup surfaces here as a
// TimeoutException instead of hanging the whole test run.
for (var round = 0; round < rounds; round++)
{
var requests = new Task<IExecutionResult>[requestsPerRound];
for (var i = 0; i < requestsPerRound; i++)
{
requests[i] = executor.ExecuteAsync(mutation, TestContext.Current.CancellationToken);
}

var results = await Task.WhenAll(requests)
.WaitAsync(TimeSpan.FromSeconds(10), TestContext.Current.CancellationToken);

foreach (var result in results)
{
Assert.Empty(result.ExpectOperationResult().Errors);
}
}
}

public class RaceMutation
{
public RacePayload DoWork() => new();
}

public class RacePayload
{
public async Task<int> Value()
{
await Task.Yield();
return 1;
}
}

public class Mutation1
{
private int _order;
Expand Down
Loading