import { describe, it } from 'node:test';
import { assert, expect } from 'chai';
import { expectJSON } from '../../../__testUtils__/expectJSON.ts';
import { expectPromise } from '../../../__testUtils__/expectPromise.ts';
import { resolveOnNextTick } from '../../../__testUtils__/resolveOnNextTick.ts';
import { spyOnMethod } from '../../../__testUtils__/spyOn.ts';
import { withAsyncUsing } from '../../../__testUtils__/withAsyncUsing.ts';
import type { PromiseOrValue } from '../../../jsutils/PromiseOrValue.ts';
import { promiseWithResolvers } from '../../../jsutils/promiseWithResolvers.ts';
import type { DocumentNode } from '../../../language/ast.ts';
import { parse } from '../../../language/parser.ts';
import {
GraphQLList,
GraphQLNonNull,
GraphQLObjectType,
} from '../../../type/definition.ts';
import { GraphQLID, GraphQLString } from '../../../type/scalars.ts';
import { GraphQLSchema } from '../../../type/schema.ts';
import { buildSchema } from '../../../utilities/buildASTSchema.ts';
import type {
LegacyInitialIncrementalExecutionResult,
LegacySubsequentIncrementalExecutionResult,
} from '../BranchingIncrementalExecutor.ts';
import { legacyExecuteIncrementally } from '../legacyExecuteIncrementally.ts';
const friendType = new GraphQLObjectType({
fields: {
id: { type: GraphQLID },
name: { type: GraphQLString },
nonNullName: { type: new GraphQLNonNull(GraphQLString) },
},
name: 'Friend',
});
const friends = [
{ name: 'Luke', id: 1 },
{ name: 'Han', id: 2 },
{ name: 'Leia', id: 3 },
];
const query = new GraphQLObjectType({
fields: {
scalarList: {
type: new GraphQLList(GraphQLString),
},
scalarListList: {
type: new GraphQLList(new GraphQLList(GraphQLString)),
},
friendList: {
type: new GraphQLList(friendType),
},
nonNullFriendList: {
type: new GraphQLList(new GraphQLNonNull(friendType)),
},
nestedObject: {
type: new GraphQLObjectType({
name: 'NestedObject',
fields: {
scalarField: {
type: GraphQLString,
},
nonNullScalarField: {
type: new GraphQLNonNull(GraphQLString),
},
nestedFriendList: { type: new GraphQLList(friendType) },
deeperNestedObject: {
type: new GraphQLObjectType({
name: 'DeeperNestedObject',
fields: {
nonNullScalarField: {
type: new GraphQLNonNull(GraphQLString),
},
deeperNestedFriendList: { type: new GraphQLList(friendType) },
},
}),
},
},
}),
},
},
name: 'Query',
});
const schema = new GraphQLSchema({ query });
const cancellationSchema = buildSchema(`
type Todo {
id: ID
items: [String]
author: User
}
type User {
id: ID
name: String
}
type Query {
todo: Todo
nonNullableTodo: Todo!
blocker: String
scalarList: [String]
slowScalarList: [String]
}
type Mutation {
foo: String
bar: String
}
type Subscription {
foo: String
}
`);
const cancelStreamSchema = buildSchema(`
type CancelStreamUser {
id: String
}
type CancelStreamTodo {
id: String
items: [String]
author: CancelStreamUser
}
type Query {
todos: [CancelStreamTodo]
}
`);
async function complete(
document: DocumentNode,
rootValue: unknown = {},
enableEarlyExecution = false,
) {
const result = await legacyExecuteIncrementally({
schema,
document,
rootValue,
enableEarlyExecution,
});
if ('initialResult' in result) {
const results: Array<
| LegacyInitialIncrementalExecutionResult
| LegacySubsequentIncrementalExecutionResult
> = [result.initialResult];
for await (const patch of result.subsequentResults) {
results.push(patch);
}
return results;
}
return result;
}
async function completeAsync(
document: DocumentNode,
numCalls: number,
rootValue: unknown = {},
) {
const result = await legacyExecuteIncrementally({
schema,
document,
rootValue,
});
assert('initialResult' in result);
const iterator = result.subsequentResults[Symbol.asyncIterator]();
const promises: Array<
PromiseOrValue<
IteratorResult<
| LegacyInitialIncrementalExecutionResult
| LegacySubsequentIncrementalExecutionResult
>
>
> = [{ done: false, value: result.initialResult }];
for (let i = 0; i < numCalls; i++) {
promises.push(iterator.next());
}
return Promise.all(promises);
}
describe('Execute: stream directive (legacy)', () => {
it('Can stream a list field', async () => {
const document = parse('{ scalarList @stream(initialCount: 1) }');
const result = await complete(document, {
scalarList: () => ['apple', 'banana', 'coconut'],
});
expectJSON(result).toDeepEqual([
{
data: {
scalarList: ['apple'],
},
hasNext: true,
},
{
incremental: [
{
items: ['banana', 'coconut'],
path: ['scalarList', 1],
},
],
hasNext: false,
},
]);
});
it('Can use default value of initialCount', async () => {
const document = parse('{ scalarList @stream }');
const result = await complete(document, {
scalarList: () => ['apple', 'banana', 'coconut'],
});
expectJSON(result).toDeepEqual([
{
data: {
scalarList: [],
},
hasNext: true,
},
{
incremental: [
{
items: ['apple', 'banana', 'coconut'],
path: ['scalarList', 0],
},
],
hasNext: false,
},
]);
});
it('Negative values of initialCount throw field errors', async () => {
const document = parse('{ scalarList @stream(initialCount: -2) }');
const result = await complete(document, {
scalarList: () => ['apple', 'banana', 'coconut'],
});
expectJSON(result).toDeepEqual({
errors: [
{
message: 'initialCount must be a positive integer',
locations: [
{
line: 1,
column: 3,
},
],
path: ['scalarList'],
},
],
data: {
scalarList: null,
},
});
});
it('Returns label from stream directive', async () => {
const document = parse(
'{ scalarList @stream(initialCount: 1, label: "scalar-stream") }',
);
const result = await complete(document, {
scalarList: () => ['apple', 'banana', 'coconut'],
});
expectJSON(result).toDeepEqual([
{
data: {
scalarList: ['apple'],
},
hasNext: true,
},
{
incremental: [
{
items: ['banana', 'coconut'],
path: ['scalarList', 1],
label: 'scalar-stream',
},
],
hasNext: false,
},
]);
});
it('Treats null stream label the same as no label', async () => {
const document = parse(
'{ scalarList @stream(initialCount: 1, label: null) }',
);
const result = await complete(document, {
scalarList: () => ['apple', 'banana', 'coconut'],
});
expectJSON(result).toDeepEqual([
{
data: {
scalarList: ['apple'],
},
hasNext: true,
},
{
incremental: [
{
items: ['banana', 'coconut'],
path: ['scalarList', 1],
},
],
hasNext: false,
},
]);
});
it('Can disable @stream using if argument', async () => {
const document = parse(
'{ scalarList @stream(initialCount: 0, if: false) }',
);
const result = await complete(document, {
scalarList: () => ['apple', 'banana', 'coconut'],
});
expectJSON(result).toDeepEqual({
data: { scalarList: ['apple', 'banana', 'coconut'] },
});
});
it('Does not disable stream with null if argument', async () => {
const document = parse(
'query ($shouldStream: Boolean) { scalarList @stream(initialCount: 2, if: $shouldStream) }',
);
const result = await complete(document, {
scalarList: () => ['apple', 'banana', 'coconut'],
});
expectJSON(result).toDeepEqual([
{
data: { scalarList: ['apple', 'banana'] },
hasNext: true,
},
{
incremental: [{ items: ['coconut'], path: ['scalarList', 2] }],
hasNext: false,
},
]);
});
it('Can stream multi-dimensional lists', async () => {
const document = parse('{ scalarListList @stream(initialCount: 1) }');
const result = await complete(document, {
scalarListList: () => [
['apple', 'apple', 'apple'],
['banana', 'banana', 'banana'],
['coconut', 'coconut', 'coconut'],
],
});
expectJSON(result).toDeepEqual([
{
data: {
scalarListList: [['apple', 'apple', 'apple']],
},
hasNext: true,
},
{
incremental: [
{
items: [
['banana', 'banana', 'banana'],
['coconut', 'coconut', 'coconut'],
],
path: ['scalarListList', 1],
},
],
hasNext: false,
},
]);
});
it('Can stream a field that returns a list of promises', async () => {
const document = parse(`
query {
friendList @stream(initialCount: 2) {
name
id
}
}
`);
const result = await complete(document, {
friendList: () => friends.map((f) => Promise.resolve(f)),
});
expectJSON(result).toDeepEqual([
{
data: {
friendList: [
{
name: 'Luke',
id: '1',
},
{
name: 'Han',
id: '2',
},
],
},
hasNext: true,
},
{
incremental: [
{
items: [
{
name: 'Leia',
id: '3',
},
],
path: ['friendList', 2],
},
],
hasNext: false,
},
]);
});
it('Can stream in correct order with lists of promises', async () => {
const document = parse(`
query {
friendList @stream(initialCount: 0) {
name
id
}
}
`);
const result = await complete(document, {
friendList: () => friends.map((f) => Promise.resolve(f)),
});
expectJSON(result).toDeepEqual([
{
data: {
friendList: [],
},
hasNext: true,
},
{
incremental: [
{
items: [{ name: 'Luke', id: '1' }],
path: ['friendList', 0],
},
],
hasNext: true,
},
{
incremental: [
{
items: [{ name: 'Han', id: '2' }],
path: ['friendList', 1],
},
],
hasNext: true,
},
{
incremental: [
{
items: [{ name: 'Leia', id: '3' }],
path: ['friendList', 2],
},
],
hasNext: false,
},
]);
});
it('Does not execute early if not specified', async () => {
const document = parse(`
query {
friendList @stream(initialCount: 0) {
id
}
}
`);
const order: Array<number> = [];
const result = await complete(document, {
friendList: () =>
friends.map((f, i) => ({
id: async () => {
const slowness = 3 - i;
for (let j = 0; j < slowness; j++) {
await resolveOnNextTick();
}
order.push(i);
return f.id;
},
})),
});
expectJSON(result).toDeepEqual([
{
data: {
friendList: [],
},
hasNext: true,
},
{
incremental: [
{
items: [{ id: '1' }],
path: ['friendList', 0],
},
],
hasNext: true,
},
{
incremental: [
{
items: [{ id: '2' }],
path: ['friendList', 1],
},
],
hasNext: true,
},
{
incremental: [
{
items: [{ id: '3' }],
path: ['friendList', 2],
},
],
hasNext: false,
},
]);
expect(order).to.deep.equal([0, 1, 2]);
});
it('Executes early if specified', async () => {
const document = parse(`
query {
friendList @stream(initialCount: 0) {
id
}
}
`);
const order: Array<number> = [];
const result = await complete(
document,
{
friendList: () =>
friends.map((f, i) => ({
id: async () => {
const slowness = 3 - i;
for (let j = 0; j < slowness; j++) {
await resolveOnNextTick();
}
order.push(i);
return f.id;
},
})),
},
true,
);
expectJSON(result).toDeepEqual([
{
data: {
friendList: [],
},
hasNext: true,
},
{
incremental: [
{
items: [{ id: '1' }, { id: '2' }, { id: '3' }],
path: ['friendList', 0],
},
],
hasNext: false,
},
]);
expect(order).to.deep.equal([2, 1, 0]);
});
it('Can stream a field that returns a list with nested promises', async () => {
const document = parse(`
query {
friendList @stream(initialCount: 2) {
name
id
}
}
`);
const result = await complete(document, {
friendList: () =>
friends.map((f) => ({
name: Promise.resolve(f.name),
id: Promise.resolve(f.id),
})),
});
expectJSON(result).toDeepEqual([
{
data: {
friendList: [
{
name: 'Luke',
id: '1',
},
{
name: 'Han',
id: '2',
},
],
},
hasNext: true,
},
{
incremental: [
{
items: [
{
name: 'Leia',
id: '3',
},
],
path: ['friendList', 2],
},
],
hasNext: false,
},
]);
});
it('Handles rejections in a field that returns a list of promises before initialCount is reached', async () => {
const document = parse(`
query {
friendList @stream(initialCount: 2) {
name
id
}
}
`);
const result = await complete(document, {
friendList: () =>
friends.map((f, i) => {
if (i === 1) {
return Promise.reject(new Error('bad'));
}
return Promise.resolve(f);
}),
});
expectJSON(result).toDeepEqual([
{
errors: [
{
message: 'bad',
locations: [{ line: 3, column: 9 }],
path: ['friendList', 1],
},
],
data: {
friendList: [{ name: 'Luke', id: '1' }, null],
},
hasNext: true,
},
{
incremental: [
{
items: [{ name: 'Leia', id: '3' }],
path: ['friendList', 2],
},
],
hasNext: false,
},
]);
});
it('Handles rejections in a field that returns a list of promises after initialCount is reached', async () => {
const document = parse(`
query {
friendList @stream(initialCount: 1) {
name
id
}
}
`);
const result = await complete(document, {
friendList: () =>
friends.map((f, i) => {
if (i === 1) {
return Promise.reject(new Error('bad'));
}
return Promise.resolve(f);
}),
});
expectJSON(result).toDeepEqual([
{
data: {
friendList: [{ name: 'Luke', id: '1' }],
},
hasNext: true,
},
{
incremental: [
{
items: [null],
errors: [
{
message: 'bad',
locations: [{ line: 3, column: 9 }],
path: ['friendList', 1],
},
],
path: ['friendList', 1],
},
],
hasNext: true,
},
{
incremental: [
{
items: [{ name: 'Leia', id: '3' }],
path: ['friendList', 2],
},
],
hasNext: false,
},
]);
});
it('Can stream a field that returns an async iterable', async () => {
const document = parse(`
query {
friendList @stream {
name
id
}
}
`);
const result = await complete(document, {
async *friendList() {
yield await Promise.resolve(friends[0]);
yield await Promise.resolve(friends[1]);
yield await Promise.resolve(friends[2]);
},
});
expectJSON(result).toDeepEqual([
{
data: {
friendList: [],
},
hasNext: true,
},
{
incremental: [
{
items: [{ name: 'Luke', id: '1' }],
path: ['friendList', 0],
},
],
hasNext: true,
},
{
incremental: [
{
items: [
{ name: 'Han', id: '2' },
{ name: 'Leia', id: '3' },
],
path: ['friendList', 1],
},
],
hasNext: false,
},
]);
});
it('Can stream a field that returns an async iterable, using a non-zero initialCount', async () => {
const document = parse(`
query {
friendList @stream(initialCount: 2) {
name
id
}
}
`);
const result = await complete(document, {
async *friendList() {
yield await Promise.resolve(friends[0]);
yield await Promise.resolve(friends[1]);
yield await Promise.resolve(friends[2]);
},
});
expectJSON(result).toDeepEqual([
{
data: {
friendList: [
{ name: 'Luke', id: '1' },
{ name: 'Han', id: '2' },
],
},
hasNext: true,
},
{
incremental: [
{
items: [{ name: 'Leia', id: '3' }],
path: ['friendList', 2],
},
],
hasNext: false,
},
]);
});
it('Negative values of initialCount throw field errors on a field that returns an async iterable', async () => {
const document = parse(`
query {
friendList @stream(initialCount: -2) {
name
id
}
}
`);
const result = await complete(document, {
async *friendList() {
},
});
expectJSON(result).toDeepEqual({
errors: [
{
message: 'initialCount must be a positive integer',
locations: [{ line: 3, column: 9 }],
path: ['friendList'],
},
],
data: {
friendList: null,
},
});
});
it('Does not execute early if not specified, when streaming from an async iterable', async () => {
const document = parse(`
query {
friendList @stream(initialCount: 0) {
id
}
}
`);
const order: Array<number> = [];
const slowFriend = async (n: number) => ({
id: async () => {
const slowness = (3 - n) * 10;
for (let j = 0; j < slowness; j++) {
await resolveOnNextTick();
}
order.push(n);
return friends[n].id;
},
});
const result = await complete(document, {
async *friendList() {
yield await Promise.resolve(slowFriend(0));
yield await Promise.resolve(slowFriend(1));
yield await Promise.resolve(slowFriend(2));
},
});
expectJSON(result).toDeepEqual([
{
data: {
friendList: [],
},
hasNext: true,
},
{
incremental: [
{
items: [{ id: '1' }],
path: ['friendList', 0],
},
],
hasNext: true,
},
{
incremental: [
{
items: [{ id: '2' }],
path: ['friendList', 1],
},
],
hasNext: true,
},
{
incremental: [
{
items: [{ id: '3' }],
path: ['friendList', 2],
},
],
hasNext: false,
},
]);
expect(order).to.deep.equal([0, 1, 2]);
});
it('Executes early if specified when streaming from an async iterable', async () => {
const document = parse(`
query {
friendList @stream(initialCount: 0) {
id
}
}
`);
const order: Array<number> = [];
const slowFriend = (n: number) => ({
id: async () => {
const slowness = (3 - n) * 10;
for (let j = 0; j < slowness; j++) {
await resolveOnNextTick();
}
order.push(n);
return friends[n].id;
},
});
const result = await complete(
document,
{
async *friendList() {
yield await Promise.resolve(slowFriend(0));
yield await Promise.resolve(slowFriend(1));
yield await Promise.resolve(slowFriend(2));
},
},
true,
);
expectJSON(result).toDeepEqual([
{
data: {
friendList: [],
},
hasNext: true,
},
{
incremental: [
{
items: [{ id: '1' }, { id: '2' }, { id: '3' }],
path: ['friendList', 0],
},
],
hasNext: false,
},
]);
expect(order).to.deep.equal([2, 1, 0]);
});
it('Can handle concurrent calls to .next() without waiting', async () => {
const document = parse(`
query {
friendList @stream(initialCount: 2) {
name
id
}
}
`);
const result = await completeAsync(document, 2, {
async *friendList() {
yield await Promise.resolve(friends[0]);
yield await Promise.resolve(friends[1]);
yield await Promise.resolve(friends[2]);
},
});
expectJSON(result).toDeepEqual([
{
done: false,
value: {
data: {
friendList: [
{ name: 'Luke', id: '1' },
{ name: 'Han', id: '2' },
],
},
hasNext: true,
},
},
{
done: false,
value: {
incremental: [
{
items: [{ name: 'Leia', id: '3' }],
path: ['friendList', 2],
},
],
hasNext: false,
},
},
{ done: true, value: undefined },
]);
});
it('Handles error thrown in async iterable before initialCount is reached', async () => {
const document = parse(`
query {
friendList @stream(initialCount: 2) {
name
id
}
}
`);
const result = await complete(document, {
async *friendList() {
yield await Promise.resolve(friends[0]);
throw new Error('bad');
},
});
expectJSON(result).toDeepEqual({
errors: [
{
message: 'bad',
locations: [{ line: 3, column: 9 }],
path: ['friendList'],
},
],
data: {
friendList: null,
},
});
});
it('Handles error thrown in async iterable after initialCount is reached', async () => {
const document = parse(`
query {
friendList @stream(initialCount: 1) {
name
id
}
}
`);
const result = await complete(document, {
async *friendList() {
yield await Promise.resolve(friends[0]);
throw new Error('bad');
},
});
expectJSON(result).toDeepEqual([
{
data: {
friendList: [{ name: 'Luke', id: '1' }],
},
hasNext: true,
},
{
incremental: [
{
items: null,
path: ['friendList'],
errors: [
{
message: 'bad',
locations: [{ line: 3, column: 9 }],
path: ['friendList'],
},
],
},
],
hasNext: false,
},
]);
});
it('Handles null returned in non-null list items after initialCount is reached', async () => {
const document = parse(`
query {
nonNullFriendList @stream(initialCount: 1) {
name
}
}
`);
const result = await complete(document, {
nonNullFriendList: () => [friends[0], null, friends[1]],
});
expectJSON(result).toDeepEqual([
{
data: {
nonNullFriendList: [{ name: 'Luke' }],
},
hasNext: true,
},
{
incremental: [
{
items: null,
path: ['nonNullFriendList'],
errors: [
{
message:
'Cannot return null for non-nullable field Query.nonNullFriendList.',
locations: [{ line: 3, column: 9 }],
path: ['nonNullFriendList', 1],
},
],
},
],
hasNext: false,
},
]);
});
it('Handles null returned in non-null async iterable list items after initialCount is reached', async () => {
const document = parse(`
query {
nonNullFriendList @stream(initialCount: 1) {
name
}
}
`);
const result = await complete(document, {
async *nonNullFriendList() {
try {
yield await Promise.resolve(friends[0]);
yield await Promise.resolve(null);
} finally {
throw new Error('Oops');
}
},
});
expectJSON(result).toDeepEqual([
{
data: {
nonNullFriendList: [{ name: 'Luke' }],
},
hasNext: true,
},
{
incremental: [
{
items: null,
path: ['nonNullFriendList'],
errors: [
{
message:
'Cannot return null for non-nullable field Query.nonNullFriendList.',
locations: [{ line: 3, column: 9 }],
path: ['nonNullFriendList', 1],
},
],
},
],
hasNext: false,
},
]);
});
it('Handles errors thrown by completeValue after initialCount is reached', async () => {
const document = parse(`
query {
scalarList @stream(initialCount: 1)
}
`);
const result = await complete(document, {
scalarList: () => [friends[0].name, {}],
});
expectJSON(result).toDeepEqual([
{
data: {
scalarList: ['Luke'],
},
hasNext: true,
},
{
incremental: [
{
items: [null],
errors: [
{
message: 'String cannot represent value: {}',
locations: [{ line: 3, column: 9 }],
path: ['scalarList', 1],
},
],
path: ['scalarList', 1],
},
],
hasNext: false,
},
]);
});
it('Handles async errors thrown by completeValue after initialCount is reached', async () => {
const document = parse(`
query {
friendList @stream(initialCount: 1) {
nonNullName
}
}
`);
const result = await complete(document, {
friendList: () => [
Promise.resolve({ nonNullName: friends[0].name }),
Promise.resolve({
nonNullName: () => Promise.reject(new Error('Oops')),
}),
Promise.resolve({ nonNullName: friends[1].name }),
],
});
expectJSON(result).toDeepEqual([
{
data: {
friendList: [{ nonNullName: 'Luke' }],
},
hasNext: true,
},
{
incremental: [
{
items: [null],
errors: [
{
message: 'Oops',
locations: [{ line: 4, column: 11 }],
path: ['friendList', 1, 'nonNullName'],
},
],
path: ['friendList', 1],
},
],
hasNext: true,
},
{
incremental: [
{
items: [{ nonNullName: 'Han' }],
path: ['friendList', 2],
},
],
hasNext: false,
},
]);
});
it('Handles nested async errors thrown by completeValue after initialCount is reached', async () => {
const document = parse(`
query {
friendList @stream(initialCount: 1) {
nonNullName
}
}
`);
const result = await complete(document, {
friendList: () => [
{ nonNullName: Promise.resolve(friends[0].name) },
{ nonNullName: Promise.reject(new Error('Oops')) },
{ nonNullName: Promise.resolve(friends[1].name) },
],
});
expectJSON(result).toDeepEqual([
{
data: {
friendList: [{ nonNullName: 'Luke' }],
},
hasNext: true,
},
{
incremental: [
{
items: [null],
errors: [
{
message: 'Oops',
locations: [{ line: 4, column: 11 }],
path: ['friendList', 1, 'nonNullName'],
},
],
path: ['friendList', 1],
},
],
hasNext: true,
},
{
incremental: [
{
items: [{ nonNullName: 'Han' }],
path: ['friendList', 2],
},
],
hasNext: false,
},
]);
});
it('Stops late stream item completion after item null bubbling', async () => {
const lateMetadataType = new GraphQLObjectType({
name: 'LateLegacyStreamMetadata',
fields: {
value: { type: GraphQLString },
},
});
const lateFriendType = new GraphQLObjectType({
name: 'LateLegacyStreamFriend',
fields: {
metadata: { type: lateMetadataType },
nonNullName: { type: new GraphQLNonNull(GraphQLString) },
},
});
const lateSchema = new GraphQLSchema({
query: new GraphQLObjectType({
name: 'LateLegacyStreamQuery',
fields: {
friendList: {
type: new GraphQLList(lateFriendType),
},
},
}),
});
const document = parse(`
query {
friendList @stream(initialCount: 1) {
metadata {
value
}
nonNullName
}
}
`);
const { promise: metadataPromise, resolve: resolveMetadata } =
promiseWithResolvers<{ value: () => string }>();
const execution = await legacyExecuteIncrementally({
schema: lateSchema,
document,
rootValue: {
friendList: () => [
Promise.resolve({
metadata: { value: 'ready' },
nonNullName: friends[0].name,
}),
Promise.resolve({
metadata: () => metadataPromise,
nonNullName: () => Promise.reject(new Error('Oops')),
}),
Promise.resolve({
metadata: { value: 'later' },
nonNullName: friends[1].name,
}),
],
},
});
assert('initialResult' in execution);
const results: Array<
| LegacyInitialIncrementalExecutionResult
| LegacySubsequentIncrementalExecutionResult
> = [execution.initialResult];
for await (const patch of execution.subsequentResults) {
results.push(patch);
}
expectJSON(results).toDeepEqual([
{
data: {
friendList: [{ metadata: { value: 'ready' }, nonNullName: 'Luke' }],
},
hasNext: true,
},
{
incremental: [
{
items: [null],
errors: [
{
message: 'Oops',
locations: [{ line: 7, column: 11 }],
path: ['friendList', 1, 'nonNullName'],
},
],
path: ['friendList', 1],
},
],
hasNext: true,
},
{
incremental: [
{
items: [{ metadata: { value: 'later' }, nonNullName: 'Han' }],
path: ['friendList', 2],
},
],
hasNext: false,
},
]);
const lateMetadata = {
value: () => 'late value',
};
const lateMetadataValueSpy = spyOnMethod(lateMetadata, 'value');
resolveMetadata(lateMetadata);
await resolveOnNextTick();
await resolveOnNextTick();
expect(lateMetadataValueSpy.callCount).to.equal(0);
});
it('Handles async errors thrown by completeValue after initialCount is reached for a non-nullable list', async () => {
const document = parse(`
query {
nonNullFriendList @stream(initialCount: 1) {
nonNullName
}
}
`);
const result = await complete(document, {
nonNullFriendList: () => [
Promise.resolve({ nonNullName: friends[0].name }),
Promise.resolve({
nonNullName: () => Promise.reject(new Error('Oops')),
}),
Promise.resolve({ nonNullName: friends[1].name }),
],
});
expectJSON(result).toDeepEqual([
{
data: {
nonNullFriendList: [{ nonNullName: 'Luke' }],
},
hasNext: true,
},
{
incremental: [
{
items: null,
path: ['nonNullFriendList'],
errors: [
{
message: 'Oops',
locations: [{ line: 4, column: 11 }],
path: ['nonNullFriendList', 1, 'nonNullName'],
},
],
},
],
hasNext: false,
},
]);
});
it('Handles nested async errors thrown by completeValue after initialCount is reached for a non-nullable list', async () => {
const document = parse(`
query {
nonNullFriendList @stream(initialCount: 1) {
nonNullName
}
}
`);
const result = await complete(document, {
nonNullFriendList: () => [
{ nonNullName: Promise.resolve(friends[0].name) },
{ nonNullName: Promise.reject(new Error('Oops')) },
{ nonNullName: Promise.resolve(friends[1].name) },
],
});
expectJSON(result).toDeepEqual([
{
data: {
nonNullFriendList: [{ nonNullName: 'Luke' }],
},
hasNext: true,
},
{
incremental: [
{
items: null,
path: ['nonNullFriendList'],
errors: [
{
message: 'Oops',
locations: [{ line: 4, column: 11 }],
path: ['nonNullFriendList', 1, 'nonNullName'],
},
],
},
],
hasNext: false,
},
]);
});
it('Handles async errors thrown by completeValue after initialCount is reached from async iterable', async () => {
const document = parse(`
query {
friendList @stream(initialCount: 1) {
nonNullName
}
}
`);
const result = await complete(document, {
async *friendList() {
yield await Promise.resolve({ nonNullName: friends[0].name });
yield await Promise.resolve({
nonNullName: () => Promise.reject(new Error('Oops')),
});
yield await Promise.resolve({ nonNullName: friends[1].name });
},
});
expectJSON(result).toDeepEqual([
{
data: {
friendList: [{ nonNullName: 'Luke' }],
},
hasNext: true,
},
{
incremental: [
{
items: [null],
errors: [
{
message: 'Oops',
locations: [{ line: 4, column: 11 }],
path: ['friendList', 1, 'nonNullName'],
},
],
path: ['friendList', 1],
},
],
hasNext: true,
},
{
incremental: [
{
items: [{ nonNullName: 'Han' }],
path: ['friendList', 2],
},
],
hasNext: false,
},
]);
});
it('Handles async errors thrown by completeValue after initialCount is reached from async generator for a non-nullable list', async () => {
const document = parse(`
query {
nonNullFriendList @stream(initialCount: 1) {
nonNullName
}
}
`);
const result = await complete(document, {
async *nonNullFriendList() {
yield await Promise.resolve({ nonNullName: friends[0].name });
yield await Promise.resolve({
nonNullName: () => Promise.reject(new Error('Oops')),
});
} ,
});
expectJSON(result).toDeepEqual([
{
data: {
nonNullFriendList: [{ nonNullName: 'Luke' }],
},
hasNext: true,
},
{
incremental: [
{
items: null,
path: ['nonNullFriendList'],
errors: [
{
message: 'Oops',
locations: [{ line: 4, column: 11 }],
path: ['nonNullFriendList', 1, 'nonNullName'],
},
],
},
],
hasNext: false,
},
]);
});
it('Handles async errors thrown by completeValue after initialCount is reached from async iterable for a non-nullable list when the async iterable does not provide a return method) ', async () => {
const document = parse(`
query {
nonNullFriendList @stream(initialCount: 1) {
nonNullName
}
}
`);
let count = 0;
const result = await complete(document, {
nonNullFriendList: {
[Symbol.asyncIterator]: () => ({
next: async () => {
switch (count++) {
case 0:
return Promise.resolve({
done: false,
value: { nonNullName: friends[0].name },
});
case 1:
return Promise.resolve({
done: false,
value: {
nonNullName: () => Promise.reject(new Error('Oops')),
},
});
case 2:
return Promise.resolve({
done: false,
value: { nonNullName: friends[1].name },
});
}
},
}),
},
});
expectJSON(result).toDeepEqual([
{
data: {
nonNullFriendList: [{ nonNullName: 'Luke' }],
},
hasNext: true,
},
{
incremental: [
{
items: null,
path: ['nonNullFriendList'],
errors: [
{
message: 'Oops',
locations: [{ line: 4, column: 11 }],
path: ['nonNullFriendList', 1, 'nonNullName'],
},
],
},
],
hasNext: false,
},
]);
});
it('Handles async errors thrown by completeValue after initialCount is reached from async iterable for a non-nullable list when the async iterable provides concurrent next/return methods and has a slow return ', async () => {
const document = parse(`
query {
nonNullFriendList @stream(initialCount: 1) {
nonNullName
}
}
`);
let count = 0;
let returned = false;
const result = await complete(document, {
nonNullFriendList: {
[Symbol.asyncIterator]: () => ({
next: async () => {
if (returned) {
return Promise.resolve({ done: true });
}
switch (count++) {
case 0:
return Promise.resolve({
done: false,
value: { nonNullName: friends[0].name },
});
case 1:
return Promise.resolve({
done: false,
value: {
nonNullName: () => Promise.reject(new Error('Oops')),
},
});
case 2:
return Promise.resolve({
done: false,
value: { nonNullName: friends[1].name },
});
}
},
return: async () => {
await resolveOnNextTick();
returned = true;
return { done: true };
},
}),
},
});
expectJSON(result).toDeepEqual([
{
data: {
nonNullFriendList: [{ nonNullName: 'Luke' }],
},
hasNext: true,
},
{
incremental: [
{
items: null,
path: ['nonNullFriendList'],
errors: [
{
message: 'Oops',
locations: [{ line: 4, column: 11 }],
path: ['nonNullFriendList', 1, 'nonNullName'],
},
],
},
],
hasNext: false,
},
]);
expect(returned).to.equal(true);
});
it('Filters payloads that are nulled', async () => {
const document = parse(`
query {
nestedObject {
nonNullScalarField
nestedFriendList @stream(initialCount: 0) {
name
}
}
}
`);
const result = await complete(document, {
nestedObject: {
nonNullScalarField: () => Promise.resolve(null),
async *nestedFriendList() {
yield await Promise.resolve(friends[0]);
} ,
},
});
expectJSON(result).toDeepEqual({
errors: [
{
message:
'Cannot return null for non-nullable field NestedObject.nonNullScalarField.',
locations: [{ line: 4, column: 11 }],
path: ['nestedObject', 'nonNullScalarField'],
},
],
data: {
nestedObject: null,
},
});
});
it('Filters payloads that are nulled by a later synchronous error', async () => {
const document = parse(`
query {
nestedObject {
nestedFriendList @stream(initialCount: 0) {
name
}
nonNullScalarField
}
}
`);
const result = await complete(document, {
nestedObject: {
async *nestedFriendList() {
yield await Promise.resolve(friends[0]);
} ,
nonNullScalarField: () => null,
},
});
expectJSON(result).toDeepEqual({
errors: [
{
message:
'Cannot return null for non-nullable field NestedObject.nonNullScalarField.',
locations: [{ line: 7, column: 11 }],
path: ['nestedObject', 'nonNullScalarField'],
},
],
data: {
nestedObject: null,
},
});
});
it('Does not filter payloads when null error is in a different path', async () => {
const document = parse(`
query {
otherNestedObject: nestedObject {
... @defer {
scalarField
}
}
nestedObject {
nestedFriendList @stream(initialCount: 0) {
name
}
}
}
`);
const result = await complete(document, {
nestedObject: {
scalarField: () => Promise.reject(new Error('Oops')),
async *nestedFriendList() {
yield await Promise.resolve(friends[0]);
},
},
});
expectJSON(result).toDeepEqual([
{
data: {
otherNestedObject: {},
nestedObject: { nestedFriendList: [] },
},
hasNext: true,
},
{
incremental: [
{
data: { scalarField: null },
errors: [
{
message: 'Oops',
locations: [{ line: 5, column: 13 }],
path: ['otherNestedObject', 'scalarField'],
},
],
path: ['otherNestedObject'],
},
{
items: [{ name: 'Luke' }],
path: ['nestedObject', 'nestedFriendList', 0],
},
],
hasNext: false,
},
]);
});
it('Cancels async stream items when null bubbles past the stream', async () => {
const document = parse(`
query {
nestedObject {
nestedFriendList @stream(initialCount: 0) {
name
}
nonNullScalarField
}
}
`);
const { promise: friendsStarted, resolve: resolveFriendsStarted } =
promiseWithResolvers<void>();
const { promise: nonNullPromise, resolve: resolveNonNull } =
promiseWithResolvers<string | null>();
const { promise: namePromise, resolve: resolveName } =
promiseWithResolvers<string>();
const resultPromise = legacyExecuteIncrementally({
schema,
document,
rootValue: {
nestedObject: {
nestedFriendList() {
resolveFriendsStarted();
return [{ name: namePromise }];
},
nonNullScalarField() {
return nonNullPromise;
},
},
},
enableEarlyExecution: true,
});
await friendsStarted;
await resolveOnNextTick();
resolveNonNull(null);
const result = await resultPromise;
assert(!('initialResult' in result));
expectJSON(result).toDeepEqual({
errors: [
{
message:
'Cannot return null for non-nullable field NestedObject.nonNullScalarField.',
locations: [{ line: 7, column: 11 }],
path: ['nestedObject', 'nonNullScalarField'],
},
],
data: {
nestedObject: null,
},
});
resolveName('Luke');
});
it('Filters stream payloads that are nulled in a deferred payload', async () => {
const document = parse(`
query {
nestedObject {
... @defer {
deeperNestedObject {
nonNullScalarField
deeperNestedFriendList @stream(initialCount: 0) {
name
}
}
}
}
}
`);
const result = await complete(document, {
nestedObject: {
deeperNestedObject: {
nonNullScalarField: () => Promise.resolve(null),
async *deeperNestedFriendList() {
yield await Promise.resolve(friends[0]);
} ,
},
},
});
expectJSON(result).toDeepEqual([
{
data: {
nestedObject: {},
},
hasNext: true,
},
{
incremental: [
{
data: {
deeperNestedObject: null,
},
errors: [
{
message:
'Cannot return null for non-nullable field DeeperNestedObject.nonNullScalarField.',
locations: [{ line: 6, column: 15 }],
path: [
'nestedObject',
'deeperNestedObject',
'nonNullScalarField',
],
},
],
path: ['nestedObject'],
},
],
hasNext: false,
},
]);
});
it('Filters defer payloads that are nulled in a stream response', async () => {
const document = parse(`
query {
friendList @stream(initialCount: 0) {
nonNullName
... @defer {
name
}
}
}
`);
const result = await complete(document, {
async *friendList() {
yield await Promise.resolve({
name: friends[0].name,
nonNullName: () => Promise.resolve(null),
});
},
});
expectJSON(result).toDeepEqual([
{
data: {
friendList: [],
},
hasNext: true,
},
{
incremental: [
{
items: [null],
errors: [
{
message:
'Cannot return null for non-nullable field Friend.nonNullName.',
locations: [{ line: 4, column: 9 }],
path: ['friendList', 0, 'nonNullName'],
},
],
path: ['friendList', 0],
},
],
hasNext: false,
},
]);
});
it('Returns iterator and ignores errors when stream payloads are filtered', async () => {
let requested = false;
const iterable = {
[Symbol.asyncIterator]() {
return this;
},
next() {
if (requested) {
return Promise.reject(new Error('Oops'));
}
requested = true;
const friend = friends[0];
return Promise.resolve({
done: false,
value: {
name: friend.name,
nonNullName: null,
},
});
},
return() {
return Promise.reject(new Error('Oops'));
},
};
const returnSpy = spyOnMethod(iterable, 'return');
const document = parse(`
query {
nestedObject {
... @defer {
deeperNestedObject {
nonNullScalarField
deeperNestedFriendList @stream(initialCount: 0) {
name
}
}
}
}
}
`);
const executeResult = await legacyExecuteIncrementally({
schema,
document,
rootValue: {
nestedObject: {
deeperNestedObject: {
nonNullScalarField: () => Promise.resolve(null),
deeperNestedFriendList: iterable,
},
},
},
enableEarlyExecution: true,
});
assert('initialResult' in executeResult);
const iterator = executeResult.subsequentResults[Symbol.asyncIterator]();
const result1 = executeResult.initialResult;
expectJSON(result1).toDeepEqual({
data: {
nestedObject: {},
},
hasNext: true,
});
const result2 = await iterator.next();
expectJSON(result2).toDeepEqual({
done: false,
value: {
incremental: [
{
data: {
deeperNestedObject: null,
},
errors: [
{
message:
'Cannot return null for non-nullable field DeeperNestedObject.nonNullScalarField.',
locations: [{ line: 6, column: 15 }],
path: [
'nestedObject',
'deeperNestedObject',
'nonNullScalarField',
],
},
],
path: ['nestedObject'],
},
],
hasNext: false,
},
});
const result3 = await iterator.next();
expectJSON(result3).toDeepEqual({ done: true, value: undefined });
assert(returnSpy.callCount === 1);
});
it('Handles promises returned by completeValue after initialCount is reached', async () => {
const document = parse(`
query {
friendList @stream(initialCount: 1) {
id
name
}
}
`);
const result = await complete(document, {
async *friendList() {
yield await Promise.resolve(friends[0]);
yield await Promise.resolve(friends[1]);
yield await Promise.resolve({
id: friends[2].id,
name: () => Promise.resolve(friends[2].name),
});
},
});
expectJSON(result).toDeepEqual([
{
data: {
friendList: [{ id: '1', name: 'Luke' }],
},
hasNext: true,
},
{
incremental: [
{
items: [{ id: '2', name: 'Han' }],
path: ['friendList', 1],
},
],
hasNext: true,
},
{
incremental: [
{
items: [{ id: '3', name: 'Leia' }],
path: ['friendList', 2],
},
],
hasNext: false,
},
]);
});
it('Handles overlapping deferred and non-deferred streams', async () => {
const document = parse(`
query {
nestedObject {
nestedFriendList @stream(initialCount: 0) {
id
}
}
nestedObject {
... @defer {
nestedFriendList @stream(initialCount: 0) {
id
name
... @defer {
innerName: name
}
}
}
}
}
`);
const result = await complete(document, {
nestedObject: {
async *nestedFriendList() {
yield await Promise.resolve(friends[0]);
yield await Promise.resolve(friends[1]);
},
},
});
expectJSON(result).toDeepEqual([
{
data: {
nestedObject: {
nestedFriendList: [],
},
},
hasNext: true,
},
{
incremental: [
{
data: {
nestedFriendList: [],
},
path: ['nestedObject'],
},
{
items: [{ id: '1' }],
path: ['nestedObject', 'nestedFriendList', 0],
},
],
hasNext: true,
},
{
incremental: [
{
items: [{ id: '1', name: 'Luke' }],
path: ['nestedObject', 'nestedFriendList', 0],
},
{
items: [{ id: '2' }],
path: ['nestedObject', 'nestedFriendList', 1],
},
{
data: { innerName: 'Luke' },
path: ['nestedObject', 'nestedFriendList', 0],
},
],
hasNext: true,
},
{
incremental: [
{
items: [{ id: '2', name: 'Han' }],
path: ['nestedObject', 'nestedFriendList', 1],
},
{
data: { innerName: 'Han' },
path: ['nestedObject', 'nestedFriendList', 1],
},
],
hasNext: false,
},
]);
});
it('Re-promotes a completed stream when a slower sibling defer resolves later', async () => {
const { promise: slowFieldPromise, resolve: resolveSlowField } =
promiseWithResolvers<string>();
const document = parse(`
query {
nestedObject {
... @defer {
nestedFriendList @stream { name }
}
... @defer {
scalarField
nestedFriendList @stream { name }
}
}
}
`);
const executeResult = await legacyExecuteIncrementally({
schema,
document,
rootValue: {
nestedObject: {
nestedFriendList: () => friends,
scalarField: slowFieldPromise,
},
},
});
assert('initialResult' in executeResult);
const iterator = executeResult.subsequentResults[Symbol.asyncIterator]();
const result1 = executeResult.initialResult;
expectJSON(result1).toDeepEqual({
data: {
nestedObject: {},
},
hasNext: true,
});
const result2 = await iterator.next();
expectJSON(result2).toDeepEqual({
value: {
incremental: [
{
data: {
nestedFriendList: [],
},
path: ['nestedObject'],
},
],
hasNext: true,
},
done: false,
});
resolveSlowField('slow');
const result3 = await iterator.next();
expectJSON(result3).toDeepEqual({
value: {
incremental: [
{
items: [{ name: 'Luke' }, { name: 'Han' }, { name: 'Leia' }],
path: ['nestedObject', 'nestedFriendList', 0],
},
],
hasNext: true,
},
done: false,
});
const result4 = await iterator.next();
expectJSON(result4).toDeepEqual({
value: {
incremental: [
{
data: {
nestedFriendList: [],
scalarField: 'slow',
},
path: ['nestedObject'],
},
],
hasNext: true,
},
done: false,
});
const result5 = await iterator.next();
expectJSON(result5).toDeepEqual({
value: {
incremental: [
{
items: [{ name: 'Luke' }, { name: 'Han' }, { name: 'Leia' }],
path: ['nestedObject', 'nestedFriendList', 0],
},
],
hasNext: false,
},
done: false,
});
const result6 = await iterator.next();
expectJSON(result6).toDeepEqual({
done: true,
value: undefined,
});
});
it('Returns payloads in correct order when parent deferred fragment resolves slower than stream', async () => {
const { promise: slowFieldPromise, resolve: resolveSlowField } =
promiseWithResolvers();
const document = parse(`
query {
nestedObject {
... DeferFragment @defer
}
}
fragment DeferFragment on NestedObject {
scalarField
nestedFriendList @stream(initialCount: 0) {
name
}
}
`);
const executeResult = await legacyExecuteIncrementally({
schema,
document,
rootValue: {
nestedObject: {
scalarField: () => slowFieldPromise,
async *nestedFriendList() {
yield await Promise.resolve(friends[0]);
yield await Promise.resolve(friends[1]);
},
},
},
});
assert('initialResult' in executeResult);
const iterator = executeResult.subsequentResults[Symbol.asyncIterator]();
const result1 = executeResult.initialResult;
expectJSON(result1).toDeepEqual({
data: {
nestedObject: {},
},
hasNext: true,
});
const result2Promise = iterator.next();
resolveSlowField('slow');
const result2 = await result2Promise;
expectJSON(result2).toDeepEqual({
value: {
incremental: [
{
data: { scalarField: 'slow', nestedFriendList: [] },
path: ['nestedObject'],
},
],
hasNext: true,
},
done: false,
});
const result3 = await iterator.next();
expectJSON(result3).toDeepEqual({
value: {
incremental: [
{
items: [{ name: 'Luke' }],
path: ['nestedObject', 'nestedFriendList', 0],
},
],
hasNext: true,
},
done: false,
});
const result4 = await iterator.next();
expectJSON(result4).toDeepEqual({
value: {
incremental: [
{
items: [{ name: 'Han' }],
path: ['nestedObject', 'nestedFriendList', 1],
},
],
hasNext: false,
},
done: false,
});
const result5 = await iterator.next();
expectJSON(result5).toDeepEqual({
value: undefined,
done: true,
});
});
it('Can @defer fields that are resolved after async iterable is complete', async () => {
const { promise: slowFieldPromise, resolve: resolveSlowField } =
promiseWithResolvers();
const {
promise: iterableCompletionPromise,
resolve: resolveIterableCompletion,
} = promiseWithResolvers();
const document = parse(`
query {
friendList @stream(label:"stream-label") {
...NameFragment @defer(label: "DeferName")
id
}
}
fragment NameFragment on Friend {
name
}
`);
const executeResult = await legacyExecuteIncrementally({
schema,
document,
rootValue: {
async *friendList() {
yield await Promise.resolve(friends[0]);
yield await Promise.resolve({
id: friends[1].id,
name: () => slowFieldPromise,
});
await iterableCompletionPromise;
},
},
});
assert('initialResult' in executeResult);
const iterator = executeResult.subsequentResults[Symbol.asyncIterator]();
const result1 = executeResult.initialResult;
expectJSON(result1).toDeepEqual({
data: {
friendList: [],
},
hasNext: true,
});
const result2Promise = iterator.next();
resolveIterableCompletion(null);
const result2 = await result2Promise;
expectJSON(result2).toDeepEqual({
value: {
incremental: [
{
items: [{ id: '1' }],
path: ['friendList', 0],
label: 'stream-label',
},
{
data: { name: 'Luke' },
path: ['friendList', 0],
label: 'DeferName',
},
],
hasNext: true,
},
done: false,
});
const result3Promise = iterator.next();
resolveSlowField('Han');
const result3 = await result3Promise;
expectJSON(result3).toDeepEqual({
value: {
incremental: [
{
items: [{ id: '2' }],
path: ['friendList', 1],
label: 'stream-label',
},
],
hasNext: true,
},
done: false,
});
const result4 = await iterator.next();
expectJSON(result4).toDeepEqual({
value: {
incremental: [
{
data: { name: 'Han' },
path: ['friendList', 1],
label: 'DeferName',
},
],
hasNext: false,
},
done: false,
});
const result5 = await iterator.next();
expectJSON(result5).toDeepEqual({
value: undefined,
done: true,
});
});
it('Can @defer fields that are resolved before async iterable is complete', async () => {
const { promise: slowFieldPromise, resolve: resolveSlowField } =
promiseWithResolvers();
const {
promise: iterableCompletionPromise,
resolve: resolveIterableCompletion,
} = promiseWithResolvers();
const document = parse(`
query {
friendList @stream(initialCount: 1, label:"stream-label") {
...NameFragment @defer(label: "DeferName")
id
}
}
fragment NameFragment on Friend {
name
}
`);
const executeResult = await legacyExecuteIncrementally({
schema,
document,
rootValue: {
async *friendList() {
yield await Promise.resolve(friends[0]);
yield await Promise.resolve({
id: friends[1].id,
name: () => slowFieldPromise,
});
await iterableCompletionPromise;
},
},
});
assert('initialResult' in executeResult);
const iterator = executeResult.subsequentResults[Symbol.asyncIterator]();
const result1 = executeResult.initialResult;
expectJSON(result1).toDeepEqual({
data: {
friendList: [{ id: '1' }],
},
hasNext: true,
});
const result2Promise = iterator.next();
resolveSlowField('Han');
const result2 = await result2Promise;
expectJSON(result2).toDeepEqual({
value: {
incremental: [
{
data: { name: 'Luke' },
path: ['friendList', 0],
label: 'DeferName',
},
],
hasNext: true,
},
done: false,
});
const result3 = await iterator.next();
expectJSON(result3).toDeepEqual({
value: {
incremental: [
{
items: [{ id: '2' }],
path: ['friendList', 1],
label: 'stream-label',
},
],
hasNext: true,
},
done: false,
});
const result4 = await iterator.next();
expectJSON(result4).toDeepEqual({
value: {
incremental: [
{
data: { name: 'Han' },
path: ['friendList', 1],
label: 'DeferName',
},
],
hasNext: true,
},
done: false,
});
const result5Promise = iterator.next();
resolveIterableCompletion(null);
const result5 = await result5Promise;
expectJSON(result5).toDeepEqual({
value: {
hasNext: false,
},
done: false,
});
const result6 = await iterator.next();
expectJSON(result6).toDeepEqual({
value: undefined,
done: true,
});
});
it('Returns underlying async iterables when returned generator is returned', async () => {
const iterable = {
[Symbol.asyncIterator]() {
return this;
},
next() {
return new Promise(() => {
});
},
return() {
},
};
const returnSpy = spyOnMethod(iterable, 'return');
const document = parse(`
query {
friendList @stream(initialCount: 0) {
id
}
}
`);
const executeResult = await legacyExecuteIncrementally({
schema,
document,
rootValue: {
friendList: iterable,
},
});
assert('initialResult' in executeResult);
const iterator = executeResult.subsequentResults[Symbol.asyncIterator]();
const result1 = executeResult.initialResult;
expectJSON(result1).toDeepEqual({
data: {
friendList: [],
},
hasNext: true,
});
const result2Promise = iterator.next();
const returnPromise = iterator.return();
const result2 = await result2Promise;
expectJSON(result2).toDeepEqual({
done: true,
value: undefined,
});
await returnPromise;
assert(returnSpy.callCount === 1);
});
it('Awaits stream source async iterable return before iterator return settles', async () => {
const { promise: returnCleanup, resolve: resolveReturnCleanup } =
promiseWithResolvers<void>();
const iterable = {
[Symbol.asyncIterator]() {
return this;
},
next: () =>
new Promise(() => {
}),
return: async () => {
await returnCleanup;
return {
value: undefined,
done: true,
};
},
};
const returnSpy = spyOnMethod(iterable, 'return');
const document = parse(`
query {
friendList @stream(initialCount: 0) {
id
}
}
`);
const executeResult = await legacyExecuteIncrementally({
schema,
document,
rootValue: {
friendList: iterable,
},
});
assert('initialResult' in executeResult);
const iterator = executeResult.subsequentResults[Symbol.asyncIterator]();
const result1 = executeResult.initialResult;
expectJSON(result1).toDeepEqual({
data: {
friendList: [],
},
hasNext: true,
});
const result2Promise = iterator.next();
const returnPromise = iterator.return();
let returnSettled = false;
returnPromise.then(
() => {
returnSettled = true;
},
() => {
returnSettled = true;
},
);
await resolveOnNextTick();
expect(returnSpy.callCount).to.equal(1);
expect(returnSettled).to.equal(false);
resolveReturnCleanup();
const result2 = await result2Promise;
expectJSON(result2).toDeepEqual({
done: true,
value: undefined,
});
await returnPromise;
});
it('Can return async iterable when underlying iterable does not have a return method', async () => {
let index = 0;
const iterable = {
[Symbol.asyncIterator]: () => ({
next: () => {
const friend = friends[index++];
if (friend == null) {
return Promise.resolve({ done: true, value: undefined });
}
return Promise.resolve({ done: false, value: friend });
},
}),
};
const document = parse(`
query {
friendList @stream(initialCount: 1) {
name
id
}
}
`);
const executeResult = await legacyExecuteIncrementally({
schema,
document,
rootValue: {
friendList: iterable,
},
});
assert('initialResult' in executeResult);
const iterator = executeResult.subsequentResults[Symbol.asyncIterator]();
const result1 = executeResult.initialResult;
expectJSON(result1).toDeepEqual({
data: {
friendList: [
{
id: '1',
name: 'Luke',
},
],
},
hasNext: true,
});
await iterator.return();
const result2 = await iterator.next();
expectJSON(result2).toDeepEqual({
done: true,
value: undefined,
});
});
it('Returns underlying async iterables when returned generator is thrown', async () => {
let index = 0;
const iterable = {
[Symbol.asyncIterator]() {
return this;
},
next() {
const friend = friends[index++];
if (friend == null) {
return Promise.resolve({ done: true, value: undefined });
}
return Promise.resolve({ done: false, value: friend });
},
return() {
},
};
const returnSpy = spyOnMethod(iterable, 'return');
const document = parse(`
query {
friendList @stream(initialCount: 1) {
... @defer {
name
}
id
}
}
`);
const executeResult = await legacyExecuteIncrementally({
schema,
document,
rootValue: {
friendList: iterable,
},
});
assert('initialResult' in executeResult);
const iterator = executeResult.subsequentResults[Symbol.asyncIterator]();
const result1 = executeResult.initialResult;
expectJSON(result1).toDeepEqual({
data: {
friendList: [
{
id: '1',
},
],
},
hasNext: true,
});
await expectPromise(iterator.throw(new Error('bad'))).toRejectWith('bad');
const result2 = await iterator.next();
expectJSON(result2).toDeepEqual({
done: true,
value: undefined,
});
assert(returnSpy.callCount === 1);
});
it('Returns underlying async iterables when resource is disposed before source completion', async () => {
const iterable = {
[Symbol.asyncIterator]() {
return this;
},
next() {
return new Promise(() => {
});
},
return() {
},
};
const returnSpy = spyOnMethod(iterable, 'return');
const document = parse(`
query {
friendList @stream(initialCount: 0) {
id
}
}
`);
const executeResult = await legacyExecuteIncrementally({
schema,
document,
rootValue: {
friendList: iterable,
},
});
assert('initialResult' in executeResult);
await withAsyncUsing(
executeResult.subsequentResults[Symbol.asyncIterator](),
(iterator) => {
assert(iterator != null);
const result1 = executeResult.initialResult;
expectJSON(result1).toDeepEqual({
data: {
friendList: [],
},
hasNext: true,
});
},
);
assert(returnSpy.callCount === 1);
});
it('Does not return underlying async iterables when resource is disposed after source completion', async () => {
let index = 0;
const values = [friends[0]];
const iterable = {
[Symbol.asyncIterator]() {
return this;
},
next() {
const friend = values[index++];
if (friend == null) {
return Promise.resolve({ done: true, value: undefined });
}
return Promise.resolve({ done: false, value: friend });
},
return() {
},
};
const returnSpy = spyOnMethod(iterable, 'return');
const document = parse(`
query {
friendList @stream(initialCount: 0) {
id
}
}
`);
const executeResult = await legacyExecuteIncrementally({
schema,
document,
rootValue: {
friendList: iterable,
},
});
assert('initialResult' in executeResult);
await withAsyncUsing(
executeResult.subsequentResults[Symbol.asyncIterator](),
async (iterator) => {
const result1 = executeResult.initialResult;
expectJSON(result1).toDeepEqual({
data: {
friendList: [],
},
hasNext: true,
});
expectJSON(await iterator.next()).toDeepEqual({
done: false,
value: {
incremental: [
{
items: [{ id: '1' }],
path: ['friendList', 0],
},
],
hasNext: false,
},
});
},
);
assert(!returnSpy.callCount);
});
it('limits stream batches to the default capacity (100)', async () => {
const document = parse(`
query {
friendList @stream {
id
}
}
`);
const executeResult = await legacyExecuteIncrementally({
schema,
document,
rootValue: {
async *friendList() {
for (let i = 0; i < 101; i++) {
yield await Promise.resolve(friends[i % 3]);
}
},
},
enableEarlyExecution: true,
});
assert('initialResult' in executeResult);
const iterator = executeResult.subsequentResults[Symbol.asyncIterator]();
const result1 = executeResult.initialResult;
expectJSON(result1).toDeepEqual({
data: {
friendList: [],
},
hasNext: true,
});
await new Promise((resolve) => {
setTimeout(resolve, 5);
});
const result2 = await iterator.next();
expectJSON(result2).toDeepEqual({
done: false,
value: {
incremental: [
{
items: Array.from({ length: 100 }, (_, i) => ({
id: ((i % 3) + 1).toString(),
})),
path: ['friendList', 0],
},
],
hasNext: true,
},
});
await new Promise((resolve) => {
setTimeout(resolve, 5);
});
const result3 = await iterator.next();
expectJSON(result3).toDeepEqual({
done: false,
value: {
incremental: [
{
items: [{ id: '2' }],
path: ['friendList', 100],
},
],
hasNext: false,
},
});
const result4 = await iterator.next();
expectJSON(result4).toDeepEqual({
done: true,
value: undefined,
});
});
});
describe('Execute: stream directive (legacy cancellation)', () => {
it('should stop streamed execution when aborted', async () => {
const abortController = new AbortController();
const document = parse(`
query {
todo {
id
items @stream
}
}
`);
const resultPromise = (async () => {
const result = await legacyExecuteIncrementally({
schema: cancellationSchema,
document,
rootValue: {
todo: {
id: '1',
items: [Promise.resolve('item')],
},
},
abortSignal: abortController.signal,
});
if ('initialResult' in result) {
for await (const _patch of result.subsequentResults) {
}
}
return result;
})();
abortController.abort();
await expectPromise(resultPromise).toRejectWith(
'This operation was aborted',
);
});
it('cancels streaming when aborted during async iterator next', async () => {
const abortController = new AbortController();
const document = parse('{ scalarList @stream(initialCount: 0) }');
const { promise: nextStarted, resolve: resolveNextStarted } =
promiseWithResolvers<void>();
const { promise: nextReturned, resolve: resolveNextReturned } =
promiseWithResolvers<IteratorResult<unknown>>();
let done = false;
const asyncIterator = {
[Symbol.asyncIterator]() {
return this;
},
next() {
if (done) {
return Promise.resolve({ value: undefined, done: true });
}
done = true;
resolveNextStarted();
return nextReturned;
},
};
const result = await legacyExecuteIncrementally({
schema,
document,
rootValue: {
scalarList: () => asyncIterator,
},
abortSignal: abortController.signal,
});
assert('initialResult' in result);
const iterator = result.subsequentResults[Symbol.asyncIterator]();
const nextPromise = iterator.next();
await nextStarted;
abortController.abort();
resolveNextReturned({ value: 'value', done: false });
await expectPromise(nextPromise).toRejectWith('This operation was aborted');
});
it('waits for async stream source return cleanup before abort cancellation settles', async () => {
const abortController = new AbortController();
const document = parse('{ scalarList @stream(initialCount: 0) }');
const { promise: nextStarted, resolve: resolveNextStarted } =
promiseWithResolvers<void>();
const { promise: returnCleanup, resolve: resolveReturnCleanup } =
promiseWithResolvers<void>();
let done = false;
const asyncIterator = {
[Symbol.asyncIterator]() {
return this;
},
next() {
if (done) {
return Promise.resolve({ value: undefined, done: true });
}
done = true;
resolveNextStarted();
return new Promise<IteratorResult<unknown>>(() => {
});
},
async return() {
await returnCleanup;
return { value: undefined, done: true };
},
};
const returnSpy = spyOnMethod(asyncIterator, 'return');
const result = await legacyExecuteIncrementally({
schema,
document,
rootValue: {
scalarList: () => asyncIterator,
},
abortSignal: abortController.signal,
});
assert('initialResult' in result);
const iterator = result.subsequentResults[Symbol.asyncIterator]();
const nextPromise = iterator.next();
let nextPromiseSettled = false;
nextPromise.then(
() => {
nextPromiseSettled = true;
},
() => {
nextPromiseSettled = true;
},
);
await nextStarted;
abortController.abort();
await resolveOnNextTick();
expect(returnSpy.callCount).to.equal(1);
expect(nextPromiseSettled).to.equal(false);
resolveReturnCleanup();
await expectPromise(nextPromise).toRejectWith('This operation was aborted');
expect(nextPromiseSettled).to.equal(true);
});
it('waits for async deferred nested stream item cleanup before abort cancellation settles', async () => {
const abortController = new AbortController();
const document = parse(`
query {
todos @stream(initialCount: 0) {
id
... @defer {
author {
id
}
items @stream(initialCount: 0)
}
}
}
`);
const { promise: todosNextStarted, resolve: resolveTodosNextStarted } =
promiseWithResolvers<void>();
const { promise: idStarted, resolve: resolveIdStarted } =
promiseWithResolvers<void>();
const { promise: itemsNextStarted, resolve: resolveItemsNextStarted } =
promiseWithResolvers<void>();
const { promise: authorStarted, resolve: resolveAuthorStarted } =
promiseWithResolvers<void>();
const { promise: idPromise } = promiseWithResolvers<string>();
const { promise: authorPromise } = promiseWithResolvers<{ id: string }>();
const { promise: itemsReturnCleanup, resolve: resolveItemsReturnCleanup } =
promiseWithResolvers<void>();
const never = new Promise<IteratorResult<unknown>>(() => {
});
let yieldedTodo = false;
const itemsAsyncIterator = {
[Symbol.asyncIterator]() {
return this;
},
next() {
resolveItemsNextStarted();
return never;
},
async return() {
await itemsReturnCleanup;
return { value: undefined, done: true };
},
};
const itemsReturnSpy = spyOnMethod(itemsAsyncIterator, 'return');
const todosAsyncIterator = {
[Symbol.asyncIterator]() {
return this;
},
next() {
if (yieldedTodo) {
return never;
}
yieldedTodo = true;
resolveTodosNextStarted();
return Promise.resolve({
value: {
id() {
resolveIdStarted();
return idPromise;
},
author() {
resolveAuthorStarted();
return authorPromise;
},
items() {
return itemsAsyncIterator;
},
},
done: false,
});
},
};
const result = await legacyExecuteIncrementally({
schema: cancelStreamSchema,
document,
rootValue: {
todos: () => todosAsyncIterator,
},
abortSignal: abortController.signal,
enableEarlyExecution: true,
});
assert('initialResult' in result);
const iterator = result.subsequentResults[Symbol.asyncIterator]();
const nextPromise = iterator.next();
let nextPromiseSettled = false;
nextPromise.then(
() => {
nextPromiseSettled = true;
},
() => {
nextPromiseSettled = true;
},
);
await todosNextStarted;
await idStarted;
await authorStarted;
await itemsNextStarted;
abortController.abort();
await resolveOnNextTick();
expect(itemsReturnSpy.callCount).to.equal(1);
expect(nextPromiseSettled).to.equal(false);
resolveItemsReturnCleanup();
await expectPromise(nextPromise).toRejectWith('This operation was aborted');
expect(nextPromiseSettled).to.equal(true);
});
it('cancels streaming when aborted while item promise is pending', async () => {
const abortController = new AbortController();
const document = parse('{ scalarList @stream(initialCount: 0) }');
const { promise: itemPromise, resolve: resolveItem } =
promiseWithResolvers<string>();
const { promise: nextStarted, resolve: resolveNextStarted } =
promiseWithResolvers<void>();
let done = false;
const asyncIterator = {
[Symbol.asyncIterator]() {
return this;
},
next() {
if (done) {
return Promise.resolve({ value: undefined, done: true });
}
done = true;
resolveNextStarted();
return Promise.resolve({ value: itemPromise, done: false });
},
};
const result = await legacyExecuteIncrementally({
schema,
document,
rootValue: {
scalarList: () => asyncIterator,
},
abortSignal: abortController.signal,
});
assert('initialResult' in result);
const iterator = result.subsequentResults[Symbol.asyncIterator]();
const nextPromise = iterator.next();
await nextStarted;
abortController.abort();
resolveItem('value');
await expectPromise(nextPromise).toRejectWith('This operation was aborted');
});
it('cancels pending stream item executors with deferred work when consumer cancels', async () => {
const document = parse(`
query {
todos @stream(initialCount: 0) {
id
... @defer {
author {
id
}
}
}
}
`);
const { promise: itemPromise, resolve: resolveItem } =
promiseWithResolvers<{
id: string;
author: () => { id: string };
}>();
const { promise: secondNextStarted, resolve: resolveSecondNextStarted } =
promiseWithResolvers<void>();
let firstCall = true;
const never = new Promise<IteratorResult<unknown>>(() => {
});
const todos = {
[Symbol.asyncIterator]() {
return this;
},
next() {
if (firstCall) {
firstCall = false;
return Promise.resolve({
value: itemPromise,
done: false,
});
}
resolveSecondNextStarted();
return never;
},
return() {
return Promise.resolve({ value: undefined, done: true });
},
};
const sourceReturnSpy = spyOnMethod(todos, 'return');
const result = await legacyExecuteIncrementally({
schema: cancelStreamSchema,
document,
rootValue: { todos },
enableEarlyExecution: true,
});
assert('initialResult' in result);
const iterator = result.subsequentResults[Symbol.asyncIterator]();
const nextPromise = iterator.next();
await secondNextStarted;
await expectPromise(iterator.return()).toResolve();
await expectPromise(nextPromise).toResolve();
expect(sourceReturnSpy.callCount).to.equal(1);
const todo = {
id: 'todo',
author: () => ({ id: 'author' }),
};
const deferredAuthorSpy = spyOnMethod(todo, 'author');
resolveItem(todo);
await resolveOnNextTick();
await resolveOnNextTick();
expect(deferredAuthorSpy.callCount).to.equal(0);
});
it('stops when the stream queue is back-pressured and the consumer cancels', async () => {
const document = parse('{ scalarList @stream(initialCount: 0) }');
const { promise: reachedCapacity, resolve: resolveReachedCapacity } =
promiseWithResolvers<void>();
let count = 0;
let done = false;
const iterator = {
[Symbol.iterator]() {
return this;
},
next() {
if (done) {
return { value: undefined, done: true };
}
count += 1;
if (count === 100) {
resolveReachedCapacity();
}
if (count > 100) {
done = true;
}
return { value: String(count), done: false };
},
return() {
throw new Error('ignored return error');
},
};
const returnSpy = spyOnMethod(iterator, 'return');
const result = await legacyExecuteIncrementally({
schema,
document,
rootValue: {
scalarList: () => iterator,
},
enableEarlyExecution: true,
});
assert('initialResult' in result);
await reachedCapacity;
await resolveOnNextTick();
const stream = result.subsequentResults[Symbol.asyncIterator]();
await expectPromise(stream.return()).toResolve();
expect(returnSpy.callCount).to.equal(0);
});
it('cancels tasks and streams when aborted before initial execution finishes', async () => {
const abortController = new AbortController();
const document = parse(`
query {
todo {
id
items @stream(initialCount: 0)
... @defer {
author {
id
}
}
}
blocker
}
`);
const { promise: blockerPromise, resolve: resolveBlocker } =
promiseWithResolvers<string>();
const { promise: blockerStarted, resolve: resolveBlockerStarted } =
promiseWithResolvers<void>();
const { promise: itemsStarted, resolve: resolveItemsStarted } =
promiseWithResolvers<void>();
const resultPromise = legacyExecuteIncrementally({
schema: cancellationSchema,
document,
abortSignal: abortController.signal,
rootValue: {
blocker() {
resolveBlockerStarted();
return blockerPromise;
},
todo: {
id: 'todo',
items() {
resolveItemsStarted();
return ['a', 'b'];
},
author() {
return { id: 'author' };
},
},
},
});
await itemsStarted;
await blockerStarted;
abortController.abort();
await expectPromise(resultPromise).toRejectWith(
'This operation was aborted',
);
resolveBlocker('done');
});
it('cancels async stream source cleanup when aborted before initial execution finishes', async () => {
const abortController = new AbortController();
const document = parse(`
query {
todo {
id
items @stream(initialCount: 0)
}
blocker
}
`);
const { promise: blockerPromise, resolve: resolveBlocker } =
promiseWithResolvers<string>();
const { promise: blockerStarted, resolve: resolveBlockerStarted } =
promiseWithResolvers<void>();
const { promise: itemsStarted, resolve: resolveItemsStarted } =
promiseWithResolvers<void>();
const { promise: returnCleanup, resolve: resolveReturnCleanup } =
promiseWithResolvers<void>();
const never = new Promise<IteratorResult<string>>(() => {
});
const asyncIterator = {
[Symbol.asyncIterator]() {
return this;
},
next() {
return never;
},
async return() {
await returnCleanup;
return { value: undefined, done: true };
},
};
const sourceReturnSpy = spyOnMethod(asyncIterator, 'return');
const resultPromise = legacyExecuteIncrementally({
schema: cancellationSchema,
document,
abortSignal: abortController.signal,
rootValue: {
blocker() {
resolveBlockerStarted();
return blockerPromise;
},
todo: {
id: 'todo',
items() {
resolveItemsStarted();
return asyncIterator;
},
},
},
});
await itemsStarted;
await blockerStarted;
abortController.abort();
await resolveOnNextTick();
expect(sourceReturnSpy.callCount).to.equal(1);
await expectPromise(resultPromise).toRejectWith(
'This operation was aborted',
);
resolveReturnCleanup();
resolveBlocker('done');
await resolveOnNextTick();
});
it('should ignore repeated cancellation attempts during incremental execution', async () => {
const abortController = new AbortController();
const document = parse(`
query {
todo {
id
items @stream(initialCount: 0)
... @defer {
author {
id
}
}
}
}
`);
let nextPromiseResolved = false;
const { promise: nextStarted, resolve: resolveNextStarted } =
promiseWithResolvers<void>();
const { promise: nextPromise, resolve: resolveNext } =
promiseWithResolvers<IteratorResult<string>>();
const asyncIterator = {
[Symbol.asyncIterator]() {
return this;
},
next() {
if (nextPromiseResolved) {
return Promise.resolve({ value: undefined, done: true });
}
nextPromiseResolved = true;
resolveNextStarted();
return nextPromise;
},
return() {
return Promise.resolve({ value: undefined, done: true });
},
};
const streamReturnSpy = spyOnMethod(asyncIterator, 'return');
const { promise: authorStarted, resolve: resolveAuthorStarted } =
promiseWithResolvers<void>();
const { promise: authorPromise, resolve: resolveAuthor } =
promiseWithResolvers<{ id: string }>();
const rootValue = {
todo: {
id: 'todo',
items: () => asyncIterator,
author() {
resolveAuthorStarted();
return authorPromise;
},
},
};
const result = await legacyExecuteIncrementally({
schema: cancellationSchema,
document,
rootValue,
enableEarlyExecution: true,
abortSignal: abortController.signal,
});
assert('initialResult' in result);
const iterator = result.subsequentResults[Symbol.asyncIterator]();
const nextResultPromise = iterator.next();
await authorStarted;
await nextStarted;
abortController.abort();
await resolveOnNextTick();
await expectPromise(nextResultPromise).toRejectWith(
'This operation was aborted',
);
await resolveOnNextTick();
const followUp = await iterator.next();
expect(followUp.done).to.equal(true);
expect(streamReturnSpy.callCount).to.equal(1);
abortController.abort();
await resolveOnNextTick();
expect(streamReturnSpy.callCount).to.equal(1);
resolveNext({ value: 'value', done: false });
resolveAuthor({ id: 'author' });
});
it('cancels stream item executors with deferred work and nested streams', async () => {
const document = parse(`
query {
todos @stream(initialCount: 0) {
id
items @stream(initialCount: 0)
... @defer {
author {
id
}
}
}
}
`);
const { promise: idPromise, resolve: resolveId } =
promiseWithResolvers<string>();
const { promise: idStarted, resolve: resolveIdStarted } =
promiseWithResolvers<void>();
const result = await legacyExecuteIncrementally({
schema: cancelStreamSchema,
document,
enableEarlyExecution: true,
rootValue: {
todos: [
{
id() {
resolveIdStarted();
return idPromise;
},
items: ['a'],
author: { id: 'author' },
},
],
},
});
assert('initialResult' in result);
const iterator = result.subsequentResults[Symbol.asyncIterator]();
await idStarted;
await expectPromise(iterator.return()).toResolve();
await resolveOnNextTick();
resolveId('todo');
await resolveOnNextTick();
});
it('stops streaming when a pending stream item resolves after cancellation', async () => {
const document = parse('{ scalarList @stream(initialCount: 0) }');
const { promise: itemPromise, resolve: resolveItem } =
promiseWithResolvers<string>();
const { promise: nextStarted, resolve: resolveNextStarted } =
promiseWithResolvers<void>();
let done = false;
const iterator = {
[Symbol.iterator]() {
return this;
},
next() {
if (done) {
return { value: undefined, done: true };
}
done = true;
resolveNextStarted();
return { value: itemPromise, done: false };
},
};
const result = await legacyExecuteIncrementally({
schema,
document,
rootValue: {
scalarList: () => iterator,
},
});
assert('initialResult' in result);
const stream = result.subsequentResults[Symbol.asyncIterator]();
const nextPromise = stream.next();
await nextStarted;
const returnPromise = stream.return();
await resolveOnNextTick();
resolveItem('value');
await expectPromise(returnPromise).toResolve();
await expectPromise(nextPromise).toResolve();
});
});