Skip to content

Commit 771d526

Browse files
xperiandriclaudeCopilot
authored
Added Define.TaskSeqField for lazy async collection streaming support (IAsyncEnumerable) (#607)
* Added `Define.TaskSeqField` with streaming of `IAsyncEnumerable` results A field resolved from `IAsyncEnumerable<'T>` is enumerated into a list without directives, delivered as a whole with `@defer`, and streamed item by item with `@stream`. Streaming pulls the sequence lazily, cancels the enumeration when the subscriber disposes, and reports an enumeration error as a deferred error after the items already produced, so buffered items and sibling deferred streams are not lost. `StreamBatching` groups streamed items into batches of a fixed size or of a size computed from the sequence. The `preferredBatchSize` argument of `@stream` takes precedence. Azure `AsyncPageable<T>` does not expose its page size, so tests cover both a plain pageable and one that keeps the page size hint. The `graphql-transport-ws` middleware now sends deferred and streamed payloads immediately with `path` and `hasNext`, followed by a final `hasNext: false` payload, instead of after a fixed 5 second delay. It no longer casts payloads to a dictionary and no longer drops initial payload errors. `SubscriptionExecutionResult.Data` became `obj Skippable` and the record got `Path` and `HasNext`. `FSharp.Data.GraphQL.Shared` references `Microsoft.Bcl.AsyncInterfaces` for `netstandard2.0`. The Star Wars sample got a `Human.friendsStream` field. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> * Changed `BufferedStreamOptions` and stream helpers to use `voption` `BufferedStreamOptions.Interval` and `BufferedStreamOptions.PreferredBatchSize` are now `int voption`, so the batch size computed for a `Define.TaskSeqField` sequence is used without converting between `option` and `voption`. The `@stream` planning and buffering code and the stream event filtering follow. `Define.Input` still takes the default value of a `Nullable IntType` argument as `int option`. The optional callbacks of the `TestObserver` and `SuspendingAsyncEnumerable` test helpers are struct optional parameters. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * Rewrote subscription disposal tests as asynchronous `task` tests The tests checking that disposing a subscription stops the enumeration blocked the test thread with `Thread.Sleep` and `ManualResetEventSlim.Wait`. They now return `Task`, await `TaskCompletionSource` signals through the new `waitForTask` helper, which fails the test with a message on timeout, and wait with `Task.Delay`. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * Addressed Copilot review: ordered stream failures, bounded concurrency, lazy batching Fixes the three review comments on PR #598 (commit 764cc07): - An enumeration failure of a streamed TaskSeqField could overtake an earlier item that was still resolving asynchronously, because the failure was merged as an immediately-completing observable alongside still-running item resolutions. `Observable.ofAsyncEnumerableResolved` now awaits every resolution started before the failure before emitting it, so it always arrives last. - The same function bounds how many items are pulled from the source and resolved at the same time to `maxConcurrency`, a new optional parameter on `Define.TaskSeqField` (default `Environment.ProcessorCount`), so a fast or infinite source can no longer accumulate unbounded resolver work while streaming. - `StreamBatching.FromSource`'s callback ran whenever a TaskSeqField resolver was wrapped, so it also ran for ordinary and `@defer` queries. `IAsyncEnumerableFieldValue.GetPreferredBatchSize` now computes it lazily, only when `streamed` needs it: for a `@stream` query that does not itself supply `preferredBatchSize`. `Resolve.TaskSeq` now carries a `TaskSeqStreamingOptions` record (batching policy + max concurrency) instead of a bare `StreamBatchingPolicy`. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com> * Addressed second Copilot review: enumerator lifecycle, WS partial data, subscription race Fixes the five review comments on PR #598 (commit faaacb9): - `ofAsyncEnumerable`, `ofAsyncEnumerableResolved` and `AsyncEnumerableExtensions.toArrayAsync` acquired their enumerator before the try block, so a source throwing from `GetAsyncEnumerator` bypassed the failure handling and faulted the returned Task directly. For `ofAsyncEnumerableResolved` that meant `OnError` on the merged deferred observable of the whole query instead of `DeferredErrors` for just this field, which can drop sibling deferred results and the final `hasNext: false` payload. Acquisition now happens inside the try, and a shared `disposeEnumerator` helper also routes a throwing `DisposeAsync` through the same failure path, preferring an earlier enumeration failure if there was one. - `sendSubscriptionResponseOutput` discarded the partial data the executor can return alongside `SubscriptionErrors` and sent `data: null`; it now forwards both. `applyPlanExecutionResult`'s `Direct` branch dropped the execution errors the HTTP handler forwards; it now sends them too, with a warning log matching the other branches. - `addClientSubscription` subscribed before registering the subscription id, so a deferred observable completing synchronously ran its removal callback while the id was still absent; the helper then added the already-completed subscription, stranding the id permanently (a later `Subscribe` with the same id was rejected as already taken). A `SingleAssignmentDisposable` is now registered first and assigned after subscribing, so synchronous completion can find and remove it; assigning `Disposable` on an already-disposed instance disposes the assigned value too. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com> * Moved the shared TaskSeq test sources into Helpers TaskSeqFieldTests.fs and ObservableExtensionsTests.fs had each grown their own copies of the same async-enumerable test sources: ThrowingAsyncEnumerable, an "endless numbers" source recording pulls and disposal, "one item then the enumeration fails", "one item then DisposeAsync throws", the synchronous asyncItems/asyncRange sequence, and a delay helper differing only in argument order. Moved all of them into the shared, auto-opened Helpers module next to the existing SuspendingAsyncEnumerable and waitForTask, and updated both test files to use the shared versions instead. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> * Addressed third Copilot review: in-flight tracking, slot release on failure, placeholder cleanup ObservableExtensions.fs: ofAsyncEnumerableResolved kept every asynchronously resolved item's Task in a ResizeArray until the whole source ended, so a long or infinite @stream source retained one task per delivered item despite maxConcurrency. It also released a resolution's concurrency slot only after both awaiting it and emitting succeeded, so a failed resolution or an observer throwing while a result was delivered left the slot held forever; with maxConcurrency = 1 this deadlocked the enumeration. Replaced the task list with an in-flight counter plus a TaskCompletionSource signalled once enumeration has ended and every started resolution has settled, and moved the slot release into a finally so it always runs. A resolution failure now stops pulling further items and is delivered through onFailure the same way an enumeration failure is, after every resolution already started settles. GraphQLWebsocketMiddleware.fs: addClientSubscription registered its SingleAssignmentDisposable placeholder before calling Subscribe, so a stream throwing synchronously from Subscribe left the placeholder registered forever, permanently occupying the subscription id. Subscribe is now wrapped so a synchronous failure removes the placeholder (a no-op if a synchronous completion already did) before rethrowing for the existing per-message error handling to report. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com> * Fixed CI: the observer-throws regression test deadlocked on System.Reactive's own subscription teardown The previous commit's new test subscribed a raw IObserver<'T> whose OnNext throws to reproduce the concurrency-slot leak Copilot flagged. That part of the fix is correct (confirmed by an isolated repro run 300 times without System.Reactive: the semaphore slot is always released via the finally block). But through the real System.Reactive Subscribe(IObserver<'T>) call used by the operator, and confirmed with another isolated repro against the actual System.Reactive package, an observer's OnNext throwing makes Rx tear the subscription down itself: it disposes the subscription (cancelling the enumeration's token) before rethrowing. ofAsyncEnumerableResolved correctly treats that as "nobody is listening anymore" and skips both the onFailure delivery and OnCompleted, exactly as it does when disposed for any other reason. The test's assertion that OnCompleted still fires and onFailure still gets delivered was therefore wrong, and it hung for the test's full timeout on CI (Timeout waiting for OnCompleted), failing the build on all three OS runners. Replaced it with a test that doesn't depend on OnCompleted: it uses a source that signals a TaskCompletionSource from DisposeAsync, and waits (bounded) for that instead, which still proves the enumeration reaches disposal without hanging on the concurrency slot. Corrected the doc comment and RELEASE_NOTES.md, which both overstated that onFailure is delivered in this case. Also split the RELEASE_NOTES.md bullet about graphql-transport-ws discarding partial data into its two separate, more precise fixes (subscription data vs. Direct-result errors), per Copilot's fourth review. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com> * Rechecked stop conditions after acquiring a concurrency slot ofAsyncEnumerableResolved only checked failed()/cancellation at the top of the while loop. With maxConcurrency = 1 the loop parks in slots.WaitAsync while the single in-flight resolution runs; when that resolution fails, its finally releases the slot and the loop resumes straight into MoveNextAsync, so one more item is pulled and, if it resolves synchronously, emitted before the failure that already happened - contradicting the doc comment's claim that a failed resolution "stops the enumeration". Rechecked both conditions right after acquiring the slot, releasing it and stopping without pulling when either is set. Confirmed the race and the fix with an isolated fsi repro of the operator (100/100 runs emitted [2; -1] before this change, [-1] after), since the project's test suite can't be run standalone here (see the third-review commit's message on the AspNetCore-only compiler bug being addressed on struct-optional-params). Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com> * Addressed sixth Copilot review: recheck after MoveNextAsync, clarified item-error semantics ofAsyncEnumerableResolved only rechecked the stop conditions right after acquiring a concurrency slot, not after the following MoveNextAsync. A background resolution can fail while that move is still suspended; once it completes, the code went straight to resolving the item it produced, so with maxConcurrency > 1 an item pulled after a failure could still be resolved and, if synchronous, emitted before it. Factored the two checks into one `stopped ()` predicate and used it after both awaits. Confirmed the race and the fix with the same isolated fsi repro approach as the fifth review's fix (100/100 runs emitted the extra item before this change, 0/100 after). Also addressed the review's other thread: it read "a failed item resolution stops the enumeration" (from the doc comment, RELEASE_NOTES.md and the PR description) as meaning any per-item GraphQL error should end the stream, since resolveStreamedItem wraps every ResolverResult, including Error, into a plain StreamedItem. That wording described the AsyncVal computation itself throwing (a bug in the resolution plumbing, treated like a source failure), not an ordinary resolver error, which executeResolvers already turns into a normal, non-throwing ResolverResult.Error value - the same value @stream on an ordinary list turns into that item's DeferredErrors while continuing to stream the rest (see DeferredTests."Resolver list error"). Kept that behavior, since diverging from ordinary lists here would be surprising and isn't what GraphQL's per-field error semantics call for, and reworded the doc comment, RELEASE_NOTES.md and docs/type-system.md to make the distinction explicit. Added an execution-level regression test that pins the intended behavior: an item's own field error is delivered on its own path and the following items keep streaming. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com> * Addressed seventh Copilot review: addressable batch payloads, complete after Direct/RequestError GraphQLWebsocketMiddleware.fs sent a batch of streamed items (grouped by the field's batching policy or the query's preferredBatchSize) as one payload addressed by a path ending in the list of the batch's own indices, such as ["numbers", [0, 1]]. No graphql-transport-ws client can merge that into the response tree: a batch isn't addressable by any single index, only its individual items are. Added IncrementalPayloadSplitting, a small pure module that recognises such a path and splits the batch into one payload per item, addressed the same way a field that streams one item at a time already is (a one-element data array at a path ending in that item's own index), in the batch's own order. Each item's own errors are attributed by checking which item's path they start with, since every error the engine attaches to a batch already carries the full path of the specific item it came from - no cross-message state is needed. This keeps the engine's batching (still one buffered/merged event upstream) while making every item addressable on the wire. Verified the splitting logic with an isolated fsi repro of the algorithm (out-of-order batch, and a batch with one item's own field error) before adding the xUnit test. Also fixed graphql-transport-ws never sending complete after the single next of a Direct (query/mutation) or RequestError result, which the protocol requires ("Server dispatches the Complete message indicating that the execution has completed" after "at most one Next message"). A newer incremental-delivery wire format (pending/incremental/completed, matching graphql-js 17 and Apollo Client's GraphQL17Alpha9Handler) was considered and is worth adopting, but requires a per-field completion signal threaded through the engine's whole merged deferred/streamed/live observable, which touches the ~25 pre-existing, TaskSeqField-unrelated tests in DeferredTests.fs (exact payload positions and counts) that this PR otherwise leaves alone. Tracked as follow-up work on a separate branch. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com> * Addressed eighth Copilot review: WS error message, internal helper, conventions Request (validation) errors were sent over graphql-transport-ws as a Next followed by Complete, which the protocol reserves for results; a client reading it would see a successful, null-data result instead of the operation being terminated by an error. Routed RequestError through the terminal Error message instead, with no Complete after it (Direct still gets Next + Complete: it is a real, if partly erroneous, result). That required ServerMessage.Error and ServerRawPayload.ErrorMessages to carry GQLProblemDetails list instead of NameValueLookup list, so the error payload is a standard GraphQL error array as the protocol's ExecutionError requires, rather than an arbitrary object. Replaced the one other Error call site (the subscription catch-all's hand-built NameValueLookup, which had no message field and so was not a valid GraphQL error either) with GQLProblemDetails.Create. Found while wiring this up: RawServerMessageConverter.Write serialized the ErrorMessages and CustomResponse payloads without a preceding WritePropertyName ("payload"), unlike the ExecutionResult branch. Utf8JsonWriter throws when a value is written where a property name is expected, so every "error" message, and every "pong" carrying a payload, failed to serialize - nothing covered those two write paths. Fixed, and added two regression tests; verified the throw and the fix against a bare Utf8JsonWriter first. Made IncrementalPayloadSplitting internal instead of public - it was public only so the test project could reach it, which committed obj list paths and tuple results to the package's supported surface for no reason. Exposed it to the test assembly the same way Shared and Server already do (InternalsVisibleTo), rather than the transport implementation detail. Also replaced the one list-append (@) this PR introduced with a list expression, per the project's collection conventions. Added IcedTasks to the test project only (no transitive dependency for consumers) and rewrote SuspendingAsyncEnumerable.MoveNextAsync as a valueTask CE instead of manually wrapping a Task in a ValueTask, per the async conventions; verified the valueTask CE against the real IcedTasks package before wiring it in. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com> * Addressed ninth Copilot review: a failed non-null root field is an execution result, not a request error GQLResponseContent.RequestError was produced both before execution (validation, planning, variable coercion, a middleware, or the executor itself failing) and after it, whenever executeOperation's non-null root field failed and the error propagated to the root. The previous commit routed every RequestError through graphql-transport-ws's terminal Error message, which is wrong for the second case: per the spec, a response with a failed non-null root field is still an execution result (data is null, but present), not a request rejected before execution, so it must be sent as Next + Complete like any other result. The HTTP handler had the same conflation the other way: GQLResponse.RequestError omits data for both, so such a response was missing "data" entirely instead of carrying null. Represented the root failure as Direct (null, errs) instead of adding a new case: RequestError now means pre-execution only, and both transports already handle Direct correctly. Updated the six tests that asserted RequestError for this case (three in ExecutionTests.fs, one each in LazyEnumerationExceptionTests.fs and TaskSeqFieldTests.fs) to assert Direct with null data instead; the other 21 ensureRequestError call sites are genuine pre-execution validation/coercion/middleware failures and are unchanged. Documented the distinction on the two GQLResponseContent cases and reworded the middleware's comments and log messages to match. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com> * Fix ninth-round regression: inline argument coercion must stay a RequestError executeQueryOrMutation's final Error branch (added when the ninth review's fix made a failed non-null root field produce Direct(null, errs)) was also being reached by executeRootOperation's own getArgumentValues failure for inline (literal) root field arguments, since both errors flowed into the same collectFields-aggregated Result. That reclassified inline argument and input object coercion/validation failures as execution results with null data instead of RequestError, breaking InputObjectValidatorTests's "Execute handles validation of invalid inline input records with all fields" on CI. Inline argument coercion is now checked for every root field up front, mirroring Executor.eval's coerceVariables step for variables: if any root field's arguments fail to coerce, the whole request is rejected as RequestError before any resolver runs. Only a genuine resolver failure on a non-null root field now reaches the Direct(null, errs) branch. As a side effect, a mutation with an invalid literal argument on a later root field no longer executes the earlier root fields' resolvers first. Verified with a full solution build and the complete unit test suite locally (dotnet build FSharp.Data.GraphQL.slnx + dotnet test), matching CI's approach, instead of the previous partial per-project builds. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com> * Add regression coverage for a mixed success/error stream batch The tenth Copilot review claimed Execution.collectItems' chunk branch omits a failed item's index from `indicies` while still reserving its slot in `data`, so GraphQLWebsocketMiddleware.splitBatch's List.map2 would throw on a batch mixing a successful and a failed item. It does not: both arms of collectItems' `merge` prepend the item's index (`box index :: indicies`), and only the Ok arm additionally writes into `data` - so `indicies` and `data` always end up exactly `chunk.Length` long, with a failed item's slot left null. List.map2 never sees mismatched lengths. Added two regression tests pinning this shape rather than changing behavior: one drives the real engine through a @stream query with Fixed batching where item 0's own field resolution fails and item 1 succeeds, asserting the single DeferredErrors event this produces (TaskSeqFieldTests.fs); the other exercises splitBatch directly with a null data slot (IncrementalPayloadSplittingTests.fs). Both pass. Also addressed the same review's suppressed comment: the "emits each item as soon as its fields are resolved" test's ordering assertion depended on the default maxConcurrency (Environment.ProcessorCount), which could fail on a single-CPU runner; maxConcurrency is now set explicitly on that test. Verified with a full solution build and the complete unit test suite locally, matching CI. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com> * Streamline streaming execution: struct tuples, backgroundTask, renamed AsyncEnumerable module Avoid allocating reference tuples in the per-item and per-batch stream collection paths in Execution.fs by switching to struct tuples. In ObservableExtensions.fs, switch every enumeration/resolution loop from task to backgroundTask so continuations never resume on a subscriber's or caller's synchronization context, rename AsyncEnumerableExtensions to AsyncEnumerable, and replace the ExceptionDispatchInfo round-trip with ex.Reraise(). Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com> * PR review fix * Address PR review feedback Co-authored-by: xperiandri <2365592+xperiandri@users.noreply.github.com> * Cancel MoveNextAsync on resolution failure Co-authored-by: xperiandri <2365592+xperiandri@users.noreply.github.com> * Prefer resolution failure over cancellation Co-authored-by: xperiandri <2365592+xperiandri@users.noreply.github.com> * Use `CanceledIndependently` active pattern --------- Co-authored-by: Claude Opus 5 <noreply@anthropic.com> Co-authored-by: copilot-swe-agent[bot] <198982749+Copilot@users.noreply.github.com> Co-authored-by: xperiandri <2365592+xperiandri@users.noreply.github.com>
1 parent ad88ccc commit 771d526

29 files changed

Lines changed: 2238 additions & 102 deletions

‎Packages.props‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,8 @@
2121
<PackageReference Update="FsToolkit.ErrorHandling" Version="$(FsToolkitVersion)" />
2222
<PackageReference Update="FsToolkit.ErrorHandling.TaskResult" Version="$(FsToolkitVersion)" />
2323
<PackageReference Update="Giraffe" Version="7.*" />
24+
<PackageReference Update="IcedTasks" Version="0.11.*" />
25+
<PackageReference Update="Microsoft.Bcl.AsyncInterfaces" Version="$(SystemVersion)" />
2426
<PackageReference Update="Microsoft.Extensions.Http" Version="$(MicrosoftExtensionsVersion)" />
2527
<PackageReference Update="Microsoft.Extensions.Logging.Abstractions" Version="$(MicrosoftExtensionsVersion)" />
2628
<PackageReference Update="Microsoft.NETCore.Platforms" Version="$(SystemVersion)" />
@@ -67,9 +69,11 @@
6769
<PackageReference Update="xunit.runner.visualstudio" Version="3.1.4" />
6870
</ItemGroup>
6971
<ItemGroup Label="Tests and Samples">
72+
<PackageReference Update="Azure.Core" Version="1.*" />
7073
<PackageReference Update="CommandLineParser" Version="2.9.*" />
7174
<PackageReference Update="Donald" Version="10.1.0" />
7275
<PackageReference Update="EntityFramework" Version="1.*" />
76+
<PackageReference Update="FSharp.Control.TaskSeq" Version="1.*" />
7377
<PackageReference Update="FSharp.Data.TypeProviders" Version="1.*" />
7478
<PackageReference Update="GraphQL.Server.Ui.Altair" Version="8.*" />
7579
<PackageReference Update="GraphQL.Server.Ui.GraphiQL" Version="8.*" />

‎README.md‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -111,7 +111,7 @@ This boilerplate code can be easily reduced with a built-in implementation:
111111

112112
```fsharp
113113
let streamOptions =
114-
{ Interval = Some 2000; PreferredBatchSize = None }
114+
{ Interval = ValueSome 2000; PreferredBatchSize = ValueNone }
115115
let schemaConfig =
116116
SchemaConfig.DefaultWithBufferedStream(streamOptions)
117117
```

‎RELEASE_NOTES.md‎

Lines changed: 23 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -288,6 +288,29 @@
288288

289289
* **Breaking Change** Migrated to .NET 10
290290
* **Breaking Change** Made Relay `Edge` a read-only struct
291+
* **Breaking Change** `SubscriptionExecutionResult.Data` is now `obj Skippable`, and the record has new `Path` and `HasNext` fields for incremental delivery
292+
* **Breaking Change** `BufferedStreamOptions.Interval` and `BufferedStreamOptions.PreferredBatchSize` are now `int voption`
293+
* **Breaking Change** `ServerMessage.Error` and `ServerRawPayload.ErrorMessages` now carry `GQLProblemDetails list` instead of `NameValueLookup list`, so an `error` message's `payload` is a standard GraphQL error array as the `graphql-transport-ws` protocol requires
294+
* **Breaking Change** A query or mutation whose non-null root field fails during execution now produces a `Direct` (execution) result with `null` data instead of a `RequestError`, which is now only ever produced for a request rejected before execution (validation, planning, variable or inline argument coercion, a middleware, or the executor itself failing); HTTP and `graphql-transport-ws` responses for such a failure now carry `data: null` as the spec requires, instead of omitting `data` entirely
291295
* Added case-insensitive string comparison support to `ObjectListFilter`, including comparer-aware filter cases and GraphQL filter suffix handling
292296
* Improved Relay XML documentation comments
293297
* Changed query planning to throw `MalformedGQLQueryException` for invalid queries, `NotSupportedException` for unsupported type definition implementations and `InvalidOperationException` for internal planning errors instead of `System.Exception`, with messages naming the affected field, type and execution kind
298+
* Added `Define.TaskSeqField` for list fields resolved from `IAsyncEnumerable<'T>`, such as `taskSeq { }` or Azure SDK `AsyncPageable<T>`. Without directives the sequence is enumerated into a list, `@defer` delivers the whole list, and `@stream` delivers every item as soon as it is produced and its fields are resolved
299+
* Added cancellation of a streamed `Define.TaskSeqField` enumeration when the client unsubscribes, and delivery of a failure raised acquiring the sequence's enumerator, while enumerating, or disposing it, as a deferred error for the field, after every item already pulled has been resolved and delivered, so a slower item can never be overtaken by the error that follows it; an item resolution that throws stops the enumeration and is delivered the same way, while an item whose own fields fail is delivered as that item's deferred errors and streaming continues, exactly as for `@stream` on an ordinary list; a concurrency slot is never leaked even if delivering an item's result fails
300+
* Added `maxConcurrency` to `Define.TaskSeqField`, bounding how many items of a streamed sequence are pulled and resolved at the same time; defaults to `Environment.ProcessorCount`
301+
* Added `StreamBatching` to group streamed items of a `Define.TaskSeqField` into batches of a fixed size or of a size computed from the sequence, such as a page size kept with a paged SDK sequence. The `preferredBatchSize` argument of `@stream` takes precedence, and the batching function itself is evaluated lazily, only for a `@stream` query that does not supply its own `preferredBatchSize`
302+
* Added `Microsoft.Bcl.AsyncInterfaces` dependency of `FSharp.Data.GraphQL.Shared` for `netstandard2.0`
303+
* Added `Human.friendsStream` field to the Star Wars sample to demonstrate `@stream`
304+
* Fixed `graphql-transport-ws` delivery of `@defer` and `@stream` results, which are now sent as soon as they are produced with `path` and `hasNext` instead of after a fixed 5 second delay, followed by a final payload with `hasNext: false`
305+
* Fixed `graphql-transport-ws` failure on deferred and streamed results that are not objects, such as streamed list items and scalars
306+
* Fixed `graphql-transport-ws` dropping errors of the initial payload of a deferred result together with all its deferred results
307+
* Fixed `graphql-transport-ws` discarding the partial `data` of a subscription result that also had field errors, sending `null` instead
308+
* Fixed `graphql-transport-ws` discarding the field errors of a `Direct` (non-subscription) result, sending an empty error list instead
309+
* Fixed `graphql-transport-ws` stranding a subscription id forever when its deferred result completed synchronously, before it was registered
310+
* Fixed `Define.TaskSeqField` streaming retaining a task for every item already delivered until the sequence ends
311+
* Fixed `graphql-transport-ws` leaving a subscription id occupied when subscribing to its result failed synchronously
312+
* Fixed `graphql-transport-ws` addressing a batch of streamed items (grouped by `preferredBatchSize` or `StreamBatching`) with a `path` ending in the list of the batch's own indices, such as `["numbers", [0, 1]]`, which no client can merge into the response tree; a batch is now sent as one independently addressed payload per item instead, in the batch's own order
313+
* Fixed `graphql-transport-ws` never sending `complete` after the `next` of a query or mutation result, as the protocol requires
314+
* Fixed `graphql-transport-ws` sending a request error (rejected before execution: validation, planning, variable coercion, a middleware, or the executor itself failing) as a `next` result followed by `complete`, instead of the terminal `error` message the protocol requires for it; a query or mutation whose non-null root field fails during execution still gets `next` + `complete`, since it is a result, not a request error
315+
* Fixed `graphql-transport-ws` throwing while serializing an `error` message or a `pong` carrying a payload, since neither was written under the `payload` property name `Utf8JsonWriter` requires
316+
* Fixed a query or mutation whose root field has an invalid inline (literal) argument, such as a custom input object validator failing, being reported as a `Direct` result with `null` data instead of a `RequestError`; inline argument coercion is now checked for every root field before any of them execute, the same as variable coercion, so a mutation no longer executes earlier root fields before rejecting the request over a later one's invalid argument

‎docs/type-system.md‎

Lines changed: 46 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -78,6 +78,52 @@ let rec Person = Define.Object(name = "Person", fieldsFn = fun () -> [
7878

7979
As you may see, we defined Person object definition using *rec* keyword and instead of defining fields as a list and we used a lazily evaluated function instead.
8080

81+
### Defining fields backed by asynchronous sequences
82+
83+
When a list comes from an asynchronous source, such as a database cursor or a paged SDK client, use `Define.TaskSeqField`. Its resolver returns `IAsyncEnumerable<'T>`, which is what the `taskSeq { }` computation expression from [FSharp.Control.TaskSeq](https://github.com/fsprojects/FSharp.Control.TaskSeq) and C# async iterators produce.
84+
85+
```fsharp
86+
let getOrders (customerId : int) = taskSeq {
87+
for page in 0 .. 10 do
88+
let! orders = db.GetOrdersPageAsync (customerId, page)
89+
yield! orders
90+
}
91+
92+
Define.TaskSeqField("orders", ListOf Order, fun _ customer -> getOrders customer.Id)
93+
```
94+
95+
How the sequence is delivered depends on the query:
96+
97+
- Without directives the sequence is enumerated completely and returned as a regular list.
98+
- With `@defer` on a `Nullable (ListOf ...)` field the complete list is delivered in one deferred payload.
99+
- With `@stream` every item is delivered as soon as the sequence produces it and its fields are resolved. The enumeration is cancelled when the client unsubscribes.
100+
101+
Streamed items can be grouped into batches. The `preferredBatchSize` argument of `@stream`, available with `SchemaConfig.DefaultWithBufferedStream`, has priority. Otherwise the `batching` parameter of the field applies. It is either a fixed size or a function that reads the size from the source, such as the page size of a paged SDK sequence. The function is evaluated lazily: only for a `@stream` query that does not itself specify `preferredBatchSize`, so it never runs for an ordinary or `@defer` query.
102+
103+
```fsharp
104+
Define.TaskSeqField("orders", ListOf Order, (fun _ customer -> getOrders customer.Id), batching = StreamBatching.Fixed 50)
105+
106+
Define.TaskSeqField(
107+
"blobs",
108+
ListOf Blob,
109+
(fun _ container -> listBlobs container),
110+
batching = StreamBatching.FromSource (function
111+
| :? PagedSequence<BlobItem> as paged -> ValueSome paged.PageSize
112+
| _ -> ValueNone))
113+
```
114+
115+
Azure SDK `AsyncPageable<T>` does not expose its page size, because the size is only a hint passed to `AsPages`. To batch its items by pages, keep the hint in your own type, for example a subclass of `AsyncPageable<T>` or a wrapper, and read it in `StreamBatching.FromSource`.
116+
117+
With `@stream`, at most `maxConcurrency` items are pulled from the sequence and resolved at the same time; enumeration waits for one of them to complete before pulling the next, so a fast or infinite source cannot outrun resolution. It defaults to `Environment.ProcessorCount`.
118+
119+
```fsharp
120+
Define.TaskSeqField("orders", ListOf Order, (fun _ customer -> getOrders customer.Id), maxConcurrency = 4)
121+
```
122+
123+
An error raised while enumerating the source is delivered after every item already pulled has been resolved and delivered, so a slow item can never be overtaken by a failure that follows it. An item whose own fields fail is delivered as that item's deferred errors, and the following items are still streamed, exactly as for `@stream` on an ordinary list; only an exception that escapes the item's resolution, or the source itself, ends the stream.
124+
125+
Resolvers are captured as F# quotations. A `taskSeq { }` block that uses `let!` or `yield!` cannot be written inline in the resolver lambda, so define it in a separate function as shown above. Fields defined this way do not support `WithResolveMiddleware`.
126+
81127
## Defining an Interface
82128

83129
GraphQL interfaces are so called abstract types (along with unions). This means, that they can be used as part of the query, however query materialization must always be bound to some concrete Object type definition.

‎samples/star-wars-api/Schema.fs‎

Lines changed: 19 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2,7 +2,9 @@ namespace FSharp.Data.GraphQL.Samples.StarWarsApi
22

33
open System.Linq
44
open System.Text.Json.Serialization
5+
open System.Threading.Tasks
56
open Microsoft.FSharp.Reflection
7+
open FSharp.Control
68
open FSharp.Data.GraphQL
79
open FSharp.Data.GraphQL.Types
810
open FSharp.Data.GraphQL.Server.Relay
@@ -141,6 +143,17 @@ module Schema =
141143

142144
let getCharacter id = characters |> List.tryFind (matchesId id)
143145

146+
/// Produces friends one by one with a delay, which demonstrates the @stream directive.
147+
/// TaskSeq functions are used instead of a taskSeq block, because a taskSeq block compiled
148+
/// without optimizations does not resume correctly after an await.
149+
let getFriendsStream (friendIds : string list) =
150+
friendIds
151+
|> TaskSeq.ofList
152+
|> TaskSeq.chooseAsync (fun id -> task {
153+
do! Task.Delay 500
154+
return getCharacter id
155+
})
156+
144157
let EpisodeType =
145158
Define.Enum (
146159
name = "Episode",
@@ -226,6 +239,12 @@ module Schema =
226239
con
227240
)
228241
Define.Field ("appearsIn", ListOf EpisodeType, "Which movies they appear in.", (fun _ (h : Human) -> h.AppearsIn))
242+
Define.TaskSeqField (
243+
"friendsStream",
244+
ListOf CharacterType,
245+
"The friends of the human produced one by one. Request the field with @stream to receive each friend as soon as it is available.",
246+
fun _ (h : Human) -> getFriendsStream h.Friends
247+
)
229248
Define.Field ("homePlanet", Nullable StringType, "The home planet of the human, or null if unknown.", (fun _ h -> h.HomePlanet))
230249
]
231250
)

‎samples/star-wars-api/star-wars-api.fsproj‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -7,6 +7,7 @@
77
</PropertyGroup>
88

99
<ItemGroup Label="PackageReferences">
10+
<PackageReference Include="FSharp.Control.TaskSeq" />
1011
<PackageReference Include="FsToolkit.ErrorHandling.TaskResult" />
1112
<PackageReference Include="GraphQL.Server.Ui.Altair" />
1213
<PackageReference Include="GraphQL.Server.Ui.GraphiQL" />

‎src/FSharp.Data.GraphQL.Server.AspNetCore/FSharp.Data.GraphQL.Server.AspNetCore.fsproj‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -11,6 +11,10 @@
1111
<FrameworkReference Include="Microsoft.AspNetCore.App" />
1212
</ItemGroup>
1313

14+
<ItemGroup Label="InternalsVisibleTo">
15+
<InternalsVisibleTo Include="FSharp.Data.GraphQL.Tests" />
16+
</ItemGroup>
17+
1418
<ItemGroup>
1519
<Compile Include="Helpers.fs" />
1620
<Compile Include="RequestExecutionContext.fs" />

0 commit comments

Comments
 (0)