diff --git a/backbone/src/commonMain/kotlin/opensavvy/backbone/Backbone.kt b/backbone/src/commonMain/kotlin/opensavvy/backbone/Backbone.kt index 7a85e36..bf87d8f 100644 --- a/backbone/src/commonMain/kotlin/opensavvy/backbone/Backbone.kt +++ b/backbone/src/commonMain/kotlin/opensavvy/backbone/Backbone.kt @@ -1,6 +1,8 @@ package opensavvy.backbone import kotlinx.coroutines.flow.Flow +import kotlinx.coroutines.flow.emitAll +import kotlinx.coroutines.flow.flow /** * A common interface for API endpoints. @@ -28,6 +30,25 @@ interface Backbone { */ fun directRequest(ref: Ref): Flow> + /** + * Fetches the value associated with all [refs] in an external media (e.g. a remote server, a database). + * + * This function completely bypasses the [cache]: a request will be sent everytime this is called. + * To avoid sending unnecessary requests, use [request] instead. + * + * The returned flow is **short-lived**: it is closed after the request finishes. + * The updates regarding the various references are not ordered between each other. + * + * To fetch a single value, see [directRequest]. + * + * The default implementation simply calls [directRequest] sequentially. + */ + fun batchRequests(refs: Set>): Flow> = flow { + for (ref in refs) { + emitAll(directRequest(ref)) + } + } + companion object { /** * Fetches the value associated with a [ref] in an external media (e.g. a remote server, a database). -- 2.51.2 From 11648df8054e345a511d7a4a3ff85c32f07c6ea0 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Ivan=20=E2=80=9CCLOVIS=E2=80=9D=20Canet?= Date: Fri, 12 Aug 2022 21:46:01 +0200 Subject: [PATCH 2/3] feat(backbone): Implemented Cache.Batching, which batches multiple cache requests into a single backbone call --- .../kotlin/opensavvy/backbone/Cache.kt | 161 +++++++++++++++++- 1 file changed, 156 insertions(+), 5 deletions(-) diff --git a/backbone/src/commonMain/kotlin/opensavvy/backbone/Cache.kt b/backbone/src/commonMain/kotlin/opensavvy/backbone/Cache.kt index 96503ce..650230a 100644 --- a/backbone/src/commonMain/kotlin/opensavvy/backbone/Cache.kt +++ b/backbone/src/commonMain/kotlin/opensavvy/backbone/Cache.kt @@ -1,10 +1,20 @@ package opensavvy.backbone -import kotlinx.coroutines.flow.Flow -import kotlinx.coroutines.flow.emitAll -import kotlinx.coroutines.flow.flow +import kotlinx.coroutines.CompletableDeferred +import kotlinx.coroutines.CoroutineScope +import kotlinx.coroutines.channels.Channel +import kotlinx.coroutines.channels.ReceiveChannel +import kotlinx.coroutines.channels.SendChannel +import kotlinx.coroutines.flow.* +import kotlinx.coroutines.isActive +import kotlinx.coroutines.launch +import opensavvy.backbone.Cache.Batching import opensavvy.backbone.Cache.Default +import opensavvy.backbone.Data.Companion.initialData import opensavvy.backbone.Ref.Companion.directRequest +import kotlin.coroutines.CoroutineContext +import kotlin.coroutines.EmptyCoroutineContext +import kotlin.coroutines.coroutineContext /** * Stores information temporarily to avoid unneeded network requests. @@ -36,8 +46,8 @@ import opensavvy.backbone.Ref.Companion.directRequest * * Cache chaining is instantiated in the opposite order, like iterators (the last in the chain is the first checked, * and delegates to the previous one if they do not have the value). - * The first element of the chain, and therefore the one responsible for actually starting the request, is [Default]. - * Note that [Default] has a few implementation differences, it is not recommended to use it directly without chaining it under another implementation. + * The first element of the chain, and therefore the one responsible for actually starting the request, is [Default] or [Batching]. + * Note that both have a few implementation differences, it is not recommended to use them directly without chaining under another implementation. */ interface Cache { @@ -175,6 +185,9 @@ interface Cache { * between caches and the underlying network APIs (via [Backbone]), it takes a few liberties with the API: * - [get] doesn't actually return a long-lived flow, * - [updateAll] and [expireAll] don't actually do anything. + * + * The actual request is delegated to [Backbone.directRequest]. + * To call [Backbone.batchRequests] instead, see [Batching]. */ class Default : Cache { override fun get(ref: Ref): Flow> = flow { @@ -206,6 +219,144 @@ interface Cache { } } + /** + * Default cache implementation aimed to be used as the first link in a cache chain, batching requests. + * + * Much like [Default], this is not strictly a valid implementation of [Cache], and is only intended to be used as + * the first layer in a chain. + * More information is found in [Default]. + * + * The actual request is delegated to [Backbone.batchRequests]. + * To call [Backbone.directRequest] instead, see [Default]. + * + * The batching is implemented by having multiple workers collect the various requests. + * Because of this, **all references passed to this class should come from the same [Backbone] instance**. + */ + class Batching( + + /** + * The number of workers batching the requests. + * + * All requests are split among them. + * Each worker tries to batch as many requests as possible, then calls [Backbone.batchRequests] all at once. + * + * Increasing the number of workers may increase latency. + */ + workers: Int = 1, + + /** + * The context of execution for the workers. + * + * Use a context you control to be able to stop the workers. + */ + context: CoroutineContext = EmptyCoroutineContext, + ) : Cache { + + private val requests: SendChannel, CompletableDeferred>>>> + + init { + require(workers > 0) { "There must be at least 1 actor: found $workers" } + + val requests = Channel, CompletableDeferred>>>>() + this.requests = requests + + val scope = CoroutineScope(context) + repeat(workers) { + scope.launch { + worker(requests) + } + } + } + + private suspend fun worker(requests: ReceiveChannel, CompletableDeferred>>>>) { + while (coroutineContext.isActive) { + val batch = HashSet>() + + // Store the results + // We have to store lists of Deferred in case multiple requests to the same Ref happen to be in the same + // batch. + val results = HashMap, MutableList>>>>() + + run { + // Suspend until a first request arrives + val (ref, promise) = requests.receive() + batch.add(ref) + results.getOrPut(ref) { ArrayList() } + .add(promise) + } + + while (true) { + // Try to read as many requests as possible, without suspending + val (ref, promise) = requests.tryReceive() + .getOrNull() + ?: break + + batch.add(ref) + results.getOrPut(ref) { ArrayList() } + .add(promise) + } + + // Just in case, check that they all come from the same Backbone instance + // If they don't, batching them makes no sense + val backbone = batch.first().backbone + check(batch.all { it.backbone === backbone }) { "To start a request batch, all references must depend on the same Backbone instance. Found: $batch" } + + val states = HashMap, MutableStateFlow>>() + + // Tell all clients that their request is starting + for ((ref, promises) in results) { + val state = MutableStateFlow(ref.initialData) + + for (promise in promises) { + promise.complete(state) + } + + states[ref] = state + } + + // The channel is now empty, request all of them + backbone.batchRequests(batch) + .collect { + val state = states[it.ref] + ?: error("Could not find the state for the reference ${it.ref}, this means we received an update for a reference we did not ask for. There may a bug in $backbone.") + + state.value = it + } + } + } + + override fun get(ref: Ref): Flow> = flow { + val promise = CompletableDeferred>>() + + requests.send(ref to promise) + + promise.await() + .collect { + emit(it) + } + } + + override suspend fun updateAll(values: Iterable>) { + // This has no state, there is nothing to update + } + + override suspend fun expireAll(refs: Iterable>) { + // This has no state, there is nothing to expire + } + + override suspend fun expireAllRecursively(refs: Iterable>) { + // This has no state, there is nothing to expire + } + + override suspend fun expireAll() { + // This has no state, there is nothing to expire + } + + override suspend fun expireAllRecursively() { + // This has no state, there is nothing to expire + } + } + /** * Simple implementation that handles delegating to an [upstream] cache layer the functions that should be delegated. */ -- 2.51.2 From e6031888e844962d9fc495599e5dd57556d79b23 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Ivan=20=E2=80=9CCLOVIS=E2=80=9D=20Canet?= Date: Mon, 29 Aug 2022 22:03:08 +0200 Subject: [PATCH 3/3] tests(backbone): Added a unit test for batching caches --- .../opensavvy.backbone/cache/MemoryCacheTest.kt | 14 ++++++++++++++ 1 file changed, 14 insertions(+) diff --git a/backbone/src/commonTest/kotlin/opensavvy.backbone/cache/MemoryCacheTest.kt b/backbone/src/commonTest/kotlin/opensavvy.backbone/cache/MemoryCacheTest.kt index 2cc9222..1e9e25a 100644 --- a/backbone/src/commonTest/kotlin/opensavvy.backbone/cache/MemoryCacheTest.kt +++ b/backbone/src/commonTest/kotlin/opensavvy.backbone/cache/MemoryCacheTest.kt @@ -152,4 +152,18 @@ class MemoryCacheTest { job.cancel() } + + @Test + fun expiringMemoryBatchedCache() = runTest { + val job = Job() + val cache = Cache.Batching() + .cachedInMemory(coroutineContext + job) + .expireAfter(300.milliseconds, coroutineContext + job) + + testCache(cache) + testUpdateExpiration(cache) + testAutoExpiration(cache) + + job.cancel() + } }