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 @@ -129,6 +129,12 @@ await _queryExecutor.ExecuteBatchAsync(
resultBuffer,
variableSets.Length);

// Incremental batch items return response streams that are consumed after the request
// pipeline has completed. We transfer ownership of each streamed item's operation
// context into its stream so the context is only cleaned and returned to the pool once
// the stream has been fully consumed, mirroring the single operation path.
TransferStreamedContextOwnership(operationContextBuffer, resultBuffer, variableSets.Length);
Comment on lines +132 to +136

context.Result = new OperationResultBatch([.. resultBuffer.AsSpan(0, variableSets.Length)]);
}
catch (OperationCanceledException)
Expand Down Expand Up @@ -172,6 +178,24 @@ static void Initialize(
operationContexts[variableIndex] = operationContextOwner;
}

static void TransferStreamedContextOwnership(
OperationContextOwner[] operationContextBuffer,
IExecutionResult[] resultBuffer,
int length)
{
for (var i = 0; i < length; i++)
{
if (resultBuffer[i].IsStreamResult() && operationContextBuffer[i] is { } contextOwner)
{
resultBuffer[i].RegisterForCleanup(contextOwner);

// Ownership now belongs to the stream, so we drop it from the buffer to keep
// ReleaseResources from disposing (and pooling) the context a second time.
operationContextBuffer[i] = null!;
}
}
}

static void AbandonContexts(ref OperationContextOwner[]? operationContextBuffer, int length)
{
if (operationContextBuffer is not null)
Expand Down Expand Up @@ -201,7 +225,9 @@ static void ReleaseResources(

foreach (var contextOwner in contextOwners)
{
contextOwner.Dispose();
// Owners transferred to a response stream are cleared from the buffer and are
// disposed by that stream once it has been consumed, so we skip them here.
contextOwner?.Dispose();
}

contextOwners.Clear();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -159,23 +159,25 @@ private async Task ExecuteBatchIncrementalAsync(

for (var i = 0; i < length; ++i)
{
if (i == 0)
{
var branchId = parentContext.ExecutionBranchId;
await scheduler.WaitForCompletionAsync(branchId).ConfigureAwait(false);
parentContext.DeferExecutionCoordinator.EnqueueResult(parentContext.BuildResult());
results[i] = new ResponseStream(CreateStreamAndComplete, ExecutionResultKind.DeferredResult);
}
else
{
var context = operationContexts[i].OperationContext;
var branchId = context.ExecutionBranchId;
await scheduler.WaitForCompletionAsync(branchId).ConfigureAwait(false);
context.DeferExecutionCoordinator.EnqueueResult(context.BuildResult());
results[i] = new ResponseStream(CreateStream, ExecutionResultKind.DeferredResult);
}
var context = operationContexts[i].OperationContext;
var branchId = context.ExecutionBranchId;
await scheduler.WaitForCompletionAsync(branchId).ConfigureAwait(false);
context.DeferExecutionCoordinator.EnqueueResult(context.BuildResult());

// The first item's stream drives completion of the shared execution: once it has
// delivered every payload it awaits the scheduler and seals the shared arena. Every
// item reads from its own coordinator and honors its own cancellation token, so a
// payload is never enqueued into one coordinator and read from another.
results[i] = i == 0
? new ResponseStream(CreateStreamAndComplete, ExecutionResultKind.DeferredResult)
: CreateItemStream(context);
}

static ResponseStream CreateItemStream(OperationContext context)
=> new(
() => context.DeferExecutionCoordinator.ReadResultsAsync(context.RequestAborted),
ExecutionResultKind.DeferredResult);

async IAsyncEnumerable<OperationResult> CreateStreamAndComplete()
{
var requestAborted = parentContext.RequestAborted;
Expand All @@ -192,12 +194,6 @@ async IAsyncEnumerable<OperationResult> CreateStreamAndComplete()
await execution.ConfigureAwait(false);
memory.Seal();
}

IAsyncEnumerable<OperationResult> CreateStream()
{
var requestAborted = parentContext.RequestAborted;
return parentContext.DeferExecutionCoordinator.ReadResultsAsync(requestAborted);
}
}

private static void FillSchedulerWithWork(
Expand Down
118 changes: 118 additions & 0 deletions src/HotChocolate/Core/test/Execution.Tests/DeferTests.cs
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
using System.Text.Json;
using HotChocolate.AspNetCore;
using Microsoft.AspNetCore.Http;
using Microsoft.Extensions.DependencyInjection;
Expand All @@ -6,6 +7,123 @@ namespace HotChocolate.Execution;

public class DeferTests
{
[Fact]
public async Task VariableBatch_Defer_Should_Deliver_Payloads_Per_Item_When_Executed()
{
// arrange
// the timeout turns a lost-coordinator hang into a test failure instead of a block.
var executor = await DeferAndStreamTestSchema.CreateAsync();
using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(15));

var request = OperationRequestBuilder
.New()
.SetDocument(
"""
query ($id: ID!) {
person(id: $id) {
id
... @defer {
name
}
}
}
""")
.SetVariableValues(
new List<IReadOnlyDictionary<string, object?>>
{
new Dictionary<string, object?> { ["id"] = "UGVyc29uOjE=" },
new Dictionary<string, object?> { ["id"] = "UGVyc29uOjI=" }
})
.Build();

// act
await using var result = await executor.ExecuteAsync(request, cts.Token);

var summaries = new List<string>();
foreach (var item in result.ExpectOperationResultBatch().Results)
{
summaries.Add(await SummarizeStreamAsync(item.ExpectResponseStream(), cts.Token));
}

// assert
summaries.MatchInlineSnapshots(
[
"person.id=UGVyc29uOjE=; name=Pascal; hasNext=True,False",
"person.id=UGVyc29uOjI=; name=Rafi; hasNext=True,False"
]);
}

// Reads a single item's response stream to completion and projects it to a delivery-shape
// summary that pins the initial data, the deferred data, and the per-payload hasNext sequence
// without depending on the non-deterministic branch identifiers in the incremental envelope.
private static async Task<string> SummarizeStreamAsync(
ResponseStream stream,
CancellationToken cancellationToken)
{
string? personId = null;
string? name = null;
var hasNext = new List<bool>();

await foreach (var payload in stream.ReadResultsAsync().WithCancellation(cancellationToken))
{
await using var payloadCleanup = payload;
using var document = JsonDocument.Parse(payload.ToJson());
var root = document.RootElement;

hasNext.Add(root.TryGetProperty("hasNext", out var hasNextValue) && hasNextValue.GetBoolean());

if (root.TryGetProperty("data", out var data)
&& data.ValueKind is JsonValueKind.Object
&& data.TryGetProperty("person", out var person)
&& person.TryGetProperty("id", out var id))
{
personId = id.GetString();
}

// The deferred name must arrive inside an incremental entry's data (the deferred
// payload), never in the initial data. Reading only incremental[].data keeps the
// non-deterministic branch identifiers carried elsewhere in the envelope out of
// the summary.
if (root.TryGetProperty("incremental", out var incremental)
&& incremental.ValueKind is JsonValueKind.Array)
{
foreach (var entry in incremental.EnumerateArray())
{
if (entry.TryGetProperty("data", out var incrementalData)
&& TryFindName(incrementalData, out var deferredName))
{
name = deferredName;
}
}
}
}

return $"person.id={personId}; name={name}; hasNext={string.Join(",", hasNext)}";
}

private static bool TryFindName(JsonElement element, out string? name)
{
if (element.ValueKind is JsonValueKind.Object)
{
foreach (var property in element.EnumerateObject())
{
if (property.Name is "name" && property.Value.ValueKind is JsonValueKind.String)
{
name = property.Value.GetString();
return true;
}

if (TryFindName(property.Value, out name))
{
return true;
}
}
}

name = null;
return false;
}

[Fact]
public async Task InlineFragment_Defer()
{
Expand Down
Loading