Skip to content

Commit ee213e5

Browse files
authored
Merge pull request #831 from graphql-java-kickstart/bugfix/419
Fix suspend resolvers hanging when awaiting a data loader
2 parents 088ecfa + 932805e commit ee213e5

3 files changed

Lines changed: 115 additions & 1 deletion

File tree

‎src/main/kotlin/graphql/kickstart/tools/resolver/MethodFieldResolver.kt‎

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,8 @@ import graphql.schema.DataFetchingEnvironment
1515
import graphql.schema.GraphQLFieldDefinition
1616
import graphql.schema.GraphQLTypeUtil.isScalar
1717
import graphql.schema.LightDataFetcher
18+
import kotlinx.coroutines.CoroutineStart
19+
import kotlinx.coroutines.ensureActive
1820
import kotlinx.coroutines.future.future
1921
import org.apache.commons.lang3.reflect.TypeUtils
2022
import org.reactivestreams.Publisher
@@ -211,7 +213,10 @@ internal class MethodFieldResolverDataFetcher(
211213
val args = this.args.map { it(environment) }.toTypedArray()
212214

213215
return if (isSuspendFunction) {
214-
environment.coroutineScope().future(options.coroutineContextProvider.provide()) {
216+
// start undispatched so DataLoader loads are queued before graphql-java dispatches them,
217+
// which runs the block even if the context is already cancelled, hence ensureActive
218+
environment.coroutineScope().future(options.coroutineContextProvider.provide(), CoroutineStart.UNDISPATCHED) {
219+
ensureActive()
215220
invokeSuspend(source, method, args)?.transformWithGenericWrapper(options.genericWrappers) { environment }
216221
}
217222
} else {

‎src/test/kotlin/graphql/kickstart/tools/MethodFieldResolverDataFetcherTest.kt‎

Lines changed: 28 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,7 @@ import kotlinx.coroutines.channels.ReceiveChannel
2121
import org.junit.Test
2222
import org.reactivestreams.Publisher
2323
import org.reactivestreams.tck.TestEnvironment
24+
import java.util.concurrent.CancellationException
2425
import java.util.concurrent.CompletableFuture
2526
import kotlin.coroutines.coroutineContext
2627

@@ -57,6 +58,33 @@ class MethodFieldResolverDataFetcherTest {
5758
}
5859
}
5960

61+
@Test
62+
fun `data fetcher does not invoke suspend function if coroutineContext defined by options is cancelled`() {
63+
// setup
64+
val cancelledClass = CancelledClass()
65+
66+
val resolver = createFetcher("active", cancelledClass, options = cancelledClass.options)
67+
68+
// expect
69+
val future = resolver.get(createEnvironment(DataClass())) as CompletableFuture<*>
70+
assert(runCatching { future.get() }.exceptionOrNull() is CancellationException)
71+
assert(!cancelledClass.invoked)
72+
}
73+
74+
class CancelledClass : GraphQLResolver<DataClass> {
75+
var invoked = false
76+
77+
val options = SchemaParserOptions.Builder()
78+
.coroutineContext(Dispatchers.Default + Job().apply { cancel() })
79+
.build()
80+
81+
@Suppress("UNUSED_PARAMETER")
82+
suspend fun isActive(data: DataClass): Boolean {
83+
invoked = true
84+
return true
85+
}
86+
}
87+
6088
@ExperimentalCoroutinesApi
6189
@Test
6290
fun `canceling subscription Publisher also cancels underlying Kotlin coroutine channel`() {
Lines changed: 81 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,81 @@
1+
package graphql.kickstart.tools
2+
3+
import graphql.ExecutionInput
4+
import graphql.GraphQL
5+
import graphql.schema.DataFetchingEnvironment
6+
import kotlinx.coroutines.future.await
7+
import org.dataloader.DataLoaderFactory
8+
import org.dataloader.DataLoaderRegistry
9+
import org.junit.Test
10+
import java.util.concurrent.CompletableFuture
11+
import java.util.concurrent.TimeUnit
12+
13+
class SuspendFunctionDataLoaderTest {
14+
15+
private val schema = SchemaParser.newParser()
16+
.schemaString(
17+
"""
18+
type Query {
19+
user(id: Int!): User!
20+
users(ids: [Int!]!): [User!]!
21+
}
22+
23+
type User {
24+
id: Int!
25+
friend: User!
26+
}
27+
""")
28+
.resolvers(Query(), UserResolver())
29+
.build()
30+
.makeExecutableSchema()
31+
private val gql = GraphQL.newGraphQL(schema).build()
32+
33+
@Test
34+
fun `root suspend function can await a data loader`() {
35+
repeat(20) {
36+
val data = execute("{ users(ids: [1, 2, 3]) { id } }")
37+
38+
assertEquals(data, mapOf("users" to listOf(mapOf("id" to 1), mapOf("id" to 2), mapOf("id" to 3))))
39+
}
40+
}
41+
42+
@Test
43+
fun `nested suspend function can await a data loader`() {
44+
repeat(20) {
45+
val data = execute("{ a: user(id: 1) { friend { id } } b: user(id: 2) { friend { id } } }")
46+
47+
assertEquals(data, mapOf(
48+
"a" to mapOf("friend" to mapOf("id" to 2)),
49+
"b" to mapOf("friend" to mapOf("id" to 3))
50+
))
51+
}
52+
}
53+
54+
private fun execute(query: String): Any? {
55+
val userLoader = DataLoaderFactory.newDataLoader<Int, User> { ids ->
56+
CompletableFuture.supplyAsync { ids.map { User(it) } }
57+
}
58+
val registry = DataLoaderRegistry.newRegistry().register("user", userLoader).build()
59+
60+
// a resolver that awaits a load before it is dispatched hangs, so don't wait forever
61+
val result = gql.executeAsync(ExecutionInput.newExecutionInput(query).dataLoaderRegistry(registry))
62+
.get(5, TimeUnit.SECONDS)
63+
64+
assert(result.errors.isEmpty()) { result.errors.joinToString { it.message } }
65+
return result.getData<Any>()
66+
}
67+
68+
class Query : GraphQLQueryResolver {
69+
fun user(id: Int): User = User(id)
70+
71+
suspend fun users(ids: List<Int>, env: DataFetchingEnvironment): List<User> =
72+
env.getDataLoader<Int, User>("user")!!.loadMany(ids).await()
73+
}
74+
75+
class UserResolver : GraphQLResolver<User> {
76+
suspend fun friend(user: User, env: DataFetchingEnvironment): User =
77+
env.getDataLoader<Int, User>("user")!!.load(user.id + 1).await()
78+
}
79+
80+
data class User(val id: Int)
81+
}

0 commit comments

Comments
 (0)