-
Notifications
You must be signed in to change notification settings - Fork 23
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Avoid hasNext=true on the last incremental payload #588
Changes from 3 commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
Original file line number | Diff line number | Diff line change |
---|---|---|
|
@@ -2,15 +2,13 @@ package graphql.nadel.engine | |
|
||
import graphql.incremental.DelayedIncrementalPartialResult | ||
import graphql.nadel.engine.NadelIncrementalResultSupport.OutstandingJobCounter.OutstandingJobHandle | ||
import graphql.nadel.engine.util.copy | ||
import graphql.nadel.util.getLogger | ||
import graphql.normalized.ExecutableNormalizedOperation | ||
import kotlinx.coroutines.CompletableDeferred | ||
import kotlinx.coroutines.CoroutineScope | ||
import kotlinx.coroutines.Dispatchers | ||
import kotlinx.coroutines.Job | ||
import kotlinx.coroutines.SupervisorJob | ||
import kotlinx.coroutines.cancel | ||
import kotlinx.coroutines.channels.BufferOverflow | ||
import kotlinx.coroutines.channels.Channel | ||
import kotlinx.coroutines.flow.Flow | ||
|
@@ -145,21 +143,6 @@ class NadelIncrementalResultSupport internal constructor( | |
initialCompletionLock.complete(Unit) | ||
} | ||
|
||
fun close() { | ||
coroutineScope.cancel() | ||
} | ||
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. HUH, this needs to be called. I think it's a bug that it's not. i.e. we should cancel the jobs launched on the defer scope if the incremental result is cancelled e.g. on request closed early. We don't have to address that in this PR. I suspect we can edit fun resultFlow(): Flow<DelayedIncrementalPartialResult> {
return resultFlow
.onCompletion {
close()
}
} Where |
||
|
||
private fun quickCopy( | ||
subject: DelayedIncrementalPartialResult, | ||
hasNext: Boolean, | ||
): DelayedIncrementalPartialResult { | ||
return if (subject.hasNext() == hasNext) { | ||
subject | ||
} else { | ||
subject.copy(hasNext = hasNext) | ||
} | ||
} | ||
|
||
/** | ||
* Launches a job and increments the outstanding job handle. | ||
* | ||
|
Original file line number | Diff line number | Diff line change |
---|---|---|
|
@@ -379,7 +379,9 @@ abstract class NadelIntegrationTest( | |
// Note: there exists a IncrementalExecutionResult.getIncremental but that is part of the initial result | ||
assertTrue(result is IncrementalExecutionResult) | ||
|
||
// Fuck why delayed & incremental?? Shouldn't incremental == delayed? Why is there an optional synchronous incremental?? | ||
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Sorry, should've deleted that |
||
// Note: the spec allows for a list of "incremental" results which is sent as part of the initial | ||
// result (so not delivered in a delayed fashion). This var represents the incremental results that were | ||
// sent in a delayed fashion. | ||
val actualDelayedResponses = incrementalResults!! | ||
|
||
// Should only have one element that says hasNext=false, and it should be the last one | ||
|
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
I think this belongs in
NadelIncrementalResultSupport
This class just accumulates data. This is some incremental result logic.