Sitelet https://github.com/jasync-sql/jasync-sql/pull/97/files
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 @@ -14,9 +14,11 @@ import com.github.jasync.sql.db.util.onCompleteAsync
import java.util.concurrent.CompletableFuture

abstract class ConcreteConnectionBase(
val configuration: Configuration, override val creationTime: Long = System.currentTimeMillis()
val configuration: Configuration
) : ConcreteConnection {

override val creationTime: Long = System.currentTimeMillis()

override fun <A> inTransaction(f: (Connection) -> CompletableFuture<A>): CompletableFuture<A> {
return this.sendQuery("BEGIN").flatMapAsync(configuration.executionContext) {
val p = CompletableFuture<A>()
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -57,7 +57,6 @@ data class ConnectionPoolConfiguration @JvmOverloads constructor(
val username: String = "dbuser",
val password: String? = null,
val maxActiveConnections: Int = 1,
val maxConnectionTtl: Long? = null,
val maxIdleTime: Long = TimeUnit.MINUTES.toMillis(1),
val maxPendingQueries: Int = Int.MAX_VALUE,
val connectionValidationInterval: Long = 5000,
Expand All @@ -72,7 +71,8 @@ data class ConnectionPoolConfiguration @JvmOverloads constructor(
val maximumMessageSize: Int = 16777216,
val allocator: ByteBufAllocator = PooledByteBufAllocator.DEFAULT,
val applicationName: String? = null,
val interceptors: List<Supplier<QueryInterceptor>> = emptyList()
val interceptors: List<Supplier<QueryInterceptor>> = emptyList(),
val maxConnectionTtl: Long? = null

) {
init {
Expand Down Expand Up @@ -134,7 +134,6 @@ data class ConnectionPoolConfigurationBuilder @JvmOverloads constructor(
var password: String? = null,
var maxActiveConnections: Int = 1,
var maxIdleTime: Long = TimeUnit.MINUTES.toMillis(1),
var maxConnectionTtl: Long? = null,
var maxPendingQueries: Int = Int.MAX_VALUE,
var connectionValidationInterval: Long = 5000,
var connectionCreateTimeout: Long = 5000,
Expand All @@ -148,7 +147,8 @@ data class ConnectionPoolConfigurationBuilder @JvmOverloads constructor(
var maximumMessageSize: Int = 16777216,
var allocator: ByteBufAllocator = PooledByteBufAllocator.DEFAULT,
var applicationName: String? = null,
var interceptors: MutableList<Supplier<QueryInterceptor>> = mutableListOf<Supplier<QueryInterceptor>>()
var interceptors: MutableList<Supplier<QueryInterceptor>> = mutableListOf<Supplier<QueryInterceptor>>(),
var maxConnectionTtl: Long? = null
) {
fun build(): ConnectionPoolConfiguration = ConnectionPoolConfiguration(
host = host,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,6 @@ import com.github.jasync.sql.db.util.failed
import com.github.jasync.sql.db.util.map
import com.github.jasync.sql.db.util.mapTry
import com.github.jasync.sql.db.util.onComplete
import jdk.nashorn.internal.runtime.regexp.joni.Config.log
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.CoroutineStart
import kotlinx.coroutines.SupervisorJob
Expand Down Expand Up @@ -51,7 +50,7 @@ internal constructor(
configuration: PoolConfiguration,
testItemsPeriodically: Boolean,
extraTimeForTimeoutCompletion: Long = TimeUnit.SECONDS.toMillis(30)
) : AsyncObjectPool<T>, CoroutineScope {
) : AsyncObjectPool<T>, CoroutineScope {

@Suppress("unused", "RedundantVisibilityModifier")
public constructor(
Expand Down Expand Up @@ -329,17 +328,21 @@ private class ObjectPoolActor<T : PooledObject>(
availableItems.forEach {
val item = it.item
logger.trace { "test: ${item.id} available ${it.timeElapsed} ms" }
if (it.timeElapsed > configuration.maxIdle) {
logger.trace { "releasing idle item ${item.id}" }
item.destroy()
} else if (configuration.maxObjectTtl !=null && System.currentTimeMillis() - item.creationTime > configuration.maxObjectTtl) {
logger.trace { "releasing item past ttl ${item.id}" }
item.destroy()
} else {
val test = objectFactory.test(item)
inUseItems[item] = ItemInUseHolder(item.id, isInTest = true, testFuture = test)
test.mapTry { _, t ->
offerOrLog(GiveBack(item, CompletableFuture(), t, originalTime = it.time)) { "test item" }
when {
it.timeElapsed > configuration.maxIdle -> {
logger.trace { "releasing idle item ${item.id}" }
item.destroy()
}
configuration.maxObjectTtl != null && System.currentTimeMillis() - item.creationTime > configuration.maxObjectTtl -> {
logger.trace { "releasing item past ttl ${item.id}" }
item.destroy()
}
else -> {
val test = objectFactory.test(item)
inUseItems[item] = ItemInUseHolder(item.id, isInTest = true, testFuture = test)
test.mapTry { _, t ->
offerOrLog(GiveBack(item, CompletableFuture(), t, originalTime = it.time)) { "test item" }
}
}
}
}
Expand Down Expand Up @@ -452,6 +455,7 @@ private class ObjectPoolActor<T : PooledObject>(
private fun borrowFirstAvailableItem(future: CompletableFuture<T>): Boolean {
val itemHolder = availableItems.remove()
try {
validateTtl(itemHolder.item)
itemHolder.item.borrowTo(future)
return true
} catch (e: Exception) {
Expand All @@ -461,6 +465,13 @@ private class ObjectPoolActor<T : PooledObject>(
return false
}

private fun validateTtl(item: T) {
val age = System.currentTimeMillis() - item.creationTime
if (configuration.maxObjectTtl != null && age > configuration.maxObjectTtl) {
throw MaxTtlPassedException(item.id, age, configuration.maxObjectTtl)
}
}

private val totalItems: Int get() = inUseItems.size + inCreateItems.size + availableItems.size

private fun createNewItemPutInWaitQueue(message: Take<T>) {
Expand Down Expand Up @@ -494,15 +505,11 @@ private class ObjectPoolActor<T : PooledObject>(
}
}

private fun validate(a: T) {
val tried = objectFactory.validate(a)
private fun validate(item: T) {
val tried = objectFactory.validate(item)
when (tried) {
is Failure -> throw tried.exception
}
val age = System.currentTimeMillis() - a.creationTime
if (configuration.maxObjectTtl!=null && age > configuration.maxObjectTtl) {
throw MaxTtlPassedException(a, age, configuration.maxObjectTtl)
}
}
}

Expand Down
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
package com.github.jasync.sql.db.pool

class MaxTtlPassedException(obj: PooledObject, age: Long, maxTtl: Long) :
RuntimeException("Object ${obj.id} aged out of pool with age $age over maxTtl $maxTtl") {}
class MaxTtlPassedException(id: String, age: Long, maxTtl: Long) :
RuntimeException("Object $id passed max ttl with age $age over maxTtl $maxTtl") {}
Original file line number Diff line number Diff line change
Expand Up @@ -23,12 +23,12 @@ data class PoolConfiguration @JvmOverloads constructor(
val maxObjects: Int,
val maxIdle: Long,
val maxQueueSize: Int,
val maxObjectTtl: Long? = null,
val validationInterval: Long = 5000,
val createTimeout: Long = 5000,
val testTimeout: Long = 5000,
val queryTimeout: Long? = null,
val coroutineDispatcher: CoroutineDispatcher = Dispatchers.Default
val coroutineDispatcher: CoroutineDispatcher = Dispatchers.Default,
val maxObjectTtl: Long? = null
) {
companion object {
@Suppress("unused")
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -78,11 +78,11 @@ abstract class AbstractAsyncObjectPoolSpec<T : AsyncObjectPool<Widget>> {
//reset(factory) // Considered bad form, but necessary as we depend on previous state in these tests

//"takes maxObjects back"
val returns = taken.subList(0,taken.size-1).map {
val returns = taken.map {
p.giveBack(it).get()
}
assertEquals(4, returns.size)
(0..3).forEach {
assertEquals(5, returns.size)
(0..4).forEach {
assertThat(returns[it]).isEqualTo(p)
}

Expand All @@ -93,11 +93,6 @@ abstract class AbstractAsyncObjectPoolSpec<T : AsyncObjectPool<Widget>> {

//"destroy down to maxIdle widgets"
Thread.sleep(3000)
verify(exactly = 4) { factory.destroy(any()) }
// aged out widget should be destroyed on giveback
verifyExceptionInHierarchy(MaxTtlPassedException::class.java) {
p.giveBack(taken.last()).get()
}
verify(exactly = 5) { factory.destroy(any()) }
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,6 @@ import com.github.jasync.sql.db.util.FP
import com.github.jasync.sql.db.util.Try
import com.github.jasync.sql.db.verifyException
import org.assertj.core.api.Assertions.assertThat
import org.assertj.core.api.Assertions.assertThatExceptionOfType
import org.awaitility.kotlin.await
import org.awaitility.kotlin.matches
import org.awaitility.kotlin.untilCallTo
Expand Down Expand Up @@ -168,12 +167,15 @@ class ActorBasedObjectPoolTest {
}

@Test
fun `on giveback items pool should reclaim aged-out items`() {
fun `on take items pool should reclaim items pass ttl`() {
tested = ActorBasedObjectPool(factory, configuration.copy(maxObjectTtl = 50), false)
val widget = tested.take().get()
Thread.sleep(70)
assertThatExceptionOfType(ExecutionException::class.java).isThrownBy { tested.giveBack(widget).get() }.withCauseInstanceOf(MaxTtlPassedException::class.java)
assertThat(tested.availableItems).isEmpty()
tested.giveBack(widget).get()
val widget2 = tested.take().get()
assertThat(widget).isNotEqualTo(widget2)
assertThat(factory.created.size).isEqualTo(2)
assertThat(factory.destroyed[0]).isEqualTo(widget)
}

@Test
Expand Down