Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@ import com.expediagroup.graphql.dataloader.instrumentation.syncexhaustion.state.
import org.dataloader.DataLoader
import org.dataloader.instrumentation.DataLoaderInstrumentation
import org.dataloader.instrumentation.DataLoaderInstrumentationContext
import java.util.concurrent.atomic.AtomicInteger

/**
* Custom [DataLoaderInstrumentation] implementation that helps to calculate the state of [DataLoader]s in the
Expand All @@ -29,20 +30,23 @@ import org.dataloader.instrumentation.DataLoaderInstrumentationContext
class DataLoaderSyncExecutionExhaustedDataLoaderDispatcher(
private val syncExecutionExhaustedState: SyncExecutionExhaustedState
) : DataLoaderInstrumentation {
private val contextForSyncExecutionExhausted: DataLoaderInstrumentationContext<Any?> =
override fun beginLoad(
dataLoader: DataLoader<*, *>,
key: Any,
loadContext: Any?
): DataLoaderInstrumentationContext<Any?> =
object : DataLoaderInstrumentationContext<Any?> {
/**
* Counter the load was added to, cleared on completion so the load is only decreased once
*/
private var loadCounter: AtomicInteger? = null

override fun onDispatched() {
syncExecutionExhaustedState.onDataLoaderLoadDispatched()
loadCounter = syncExecutionExhaustedState.trackDataLoaderLoad()
}
override fun onCompleted(result: Any?, t: Throwable?) {
syncExecutionExhaustedState.onDataLoaderLoadCompleted()
loadCounter?.let(syncExecutionExhaustedState::onDataLoaderLoadCompleted)
loadCounter = null
}
}

override fun beginLoad(
dataLoader: DataLoader<*, *>,
key: Any,
loadContext: Any?
): DataLoaderInstrumentationContext<Any?> =
contextForSyncExecutionExhausted
}
Original file line number Diff line number Diff line change
Expand Up @@ -26,22 +26,27 @@ import java.util.concurrent.atomic.AtomicInteger
*/
class DataLoaderRegistryState {
/**
* Count [DataLoader.load] invocations
* Count [DataLoader.load] invocations that are not completed yet and were invoked
* after the last [DataLoaderRegistry.dispatchAll]
*/
private val loadCounter = AtomicInteger(0)
@Volatile
private var loadCounter = AtomicInteger(0)

/**
* Snapshot of [loadCounter] when [DataLoaderRegistry.dispatchAll] is invoked,
* then on every load complete decrease it
* Count [DataLoader.load] invocations that are not completed yet and were invoked
* before the last [DataLoaderRegistry.dispatchAll]
*/
private val onDispatchAllLoadCounter = AtomicInteger(0)
@Volatile
private var onDispatchAllLoadCounter = AtomicInteger(0)

/**
* Take snapshot of [loadCounter] when [DataLoaderRegistry.dispatchAll] is invoked
* Take snapshot of [loadCounter] when [DataLoaderRegistry.dispatchAll] is invoked,
* loads tracked by [trackDataLoaderLoad] before the snapshot will decrease [onDispatchAllLoadCounter] when completed
*/
@Synchronized
fun takeSnapshot() {
onDispatchAllLoadCounter.set(loadCounter.get())
loadCounter.set(0)
onDispatchAllLoadCounter = loadCounter
loadCounter = AtomicInteger(0)
}

/**
Expand All @@ -59,14 +64,46 @@ class DataLoaderRegistryState {
/**
* Increase [loadCounter] when [DataLoader.load] is invoked
*/
@Deprecated(
"Loads are tracked by SyncExecutionExhaustedState through DataLoaderSyncExecutionExhaustedDataLoaderDispatcher, " +
"register it with KotlinDataLoaderRegistryFactory.generate instead of calling this method directly. " +
"Will be removed in the next major version."
)
fun onDataLoaderLoadDispatched() {
loadCounter.incrementAndGet()
}

/**
* Decrease [onDispatchAllLoadCounter] when [DataLoader.load] returned [CompletableFuture] completes
*/
@Deprecated(
"Loads are tracked by SyncExecutionExhaustedState through DataLoaderSyncExecutionExhaustedDataLoaderDispatcher, " +
"register it with KotlinDataLoaderRegistryFactory.generate instead of calling this method directly. " +
"This method decreases onDispatchAllLoadCounter even for loads that were not included in the last snapshot. " +
"Will be removed in the next major version."
)
fun onDataLoaderLoadCompleted() {
onDispatchAllLoadCounter.decrementAndGet()
}

/**
* Increase [loadCounter] when [DataLoader.load] is invoked
*
* @return the counter that the load was added to, it needs to be provided to [onDataLoaderLoadCompleted]
* when the [CompletableFuture] returned by [DataLoader.load] completes
*/
@Synchronized
internal fun trackDataLoaderLoad(): AtomicInteger =
loadCounter.also(AtomicInteger::incrementAndGet)

/**
* Decrease the counter that the load was added to when [DataLoader.load] returned [CompletableFuture] completes,
* a load that completes before the next snapshot, for example a cache hit on a completed [CompletableFuture],
* will decrease [loadCounter] instead of [onDispatchAllLoadCounter]
*
* @param loadCounter the counter returned by [trackDataLoaderLoad]
*/
internal fun onDataLoaderLoadCompleted(loadCounter: AtomicInteger) {
loadCounter.decrementAndGet()
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -157,20 +157,53 @@ class SyncExecutionExhaustedState(
/**
* This is invoked right after a [DataLoader.load] was dispatched
*/
@Deprecated(
"Loads are tracked by DataLoaderSyncExecutionExhaustedDataLoaderDispatcher, register it with " +
"KotlinDataLoaderRegistryFactory.generate instead of calling this method directly. " +
"Will be removed in the next major version."
)
@Suppress("DEPRECATION")
fun onDataLoaderLoadDispatched() {
dataLoadersDispatchState.onDataLoaderLoadDispatched()
}

/**
* This is invoked right after a [DataLoader.load] was completed
*/
@Deprecated(
"Loads are tracked by DataLoaderSyncExecutionExhaustedDataLoaderDispatcher, register it with " +
"KotlinDataLoaderRegistryFactory.generate instead of calling this method directly. " +
"This method decreases the dispatched loads counter even for loads that were not included in the last dispatch. " +
"Will be removed in the next major version."
)
@Suppress("DEPRECATION")
fun onDataLoaderLoadCompleted() {
dataLoadersDispatchState.onDataLoaderLoadCompleted()
if (allSyncExecutionsExhausted()) {
dataLoaderRegistryProvider.invoke().dispatchAll()
}
}

/**
* This is invoked right after a [DataLoader.load] was dispatched
*
* @return the counter that the load was added to, it needs to be provided to [onDataLoaderLoadCompleted]
*/
internal fun trackDataLoaderLoad(): AtomicInteger =
dataLoadersDispatchState.trackDataLoaderLoad()

/**
* This is invoked right after a [DataLoader.load] was completed
*
* @param loadCounter the counter returned by [trackDataLoaderLoad] when the load was dispatched
*/
internal fun onDataLoaderLoadCompleted(loadCounter: AtomicInteger) {
dataLoadersDispatchState.onDataLoaderLoadCompleted(loadCounter)
if (allSyncExecutionsExhausted()) {
dataLoaderRegistryProvider.invoke().dispatchAll()
}
}

/**
* check if all [ExecutionInput] sharing a [GraphQLContext] exhausted their execution.
* A Synchronous Execution is considered Exhausted when all [DataFetcher]s of all paths were executed up until
Expand Down
Loading
Loading