mirror of
https://github.com/InsanusMokrassar/MicroUtils.git
synced 2026-09-29 11:54:53 +00:00
fixes in smart rw locker
This commit is contained in:
@@ -2,6 +2,10 @@
|
||||
|
||||
## 0.25.8.3
|
||||
|
||||
* `Coroutines`:
|
||||
* `SmartRWLocker`:
|
||||
* Fix of `unlockWrite`, `lockWrite` and `releaseRead` calls to pass correct number of permits
|
||||
|
||||
## 0.25.8.2
|
||||
|
||||
* `Coroutines`:
|
||||
|
||||
@@ -1,6 +1,8 @@
|
||||
package dev.inmo.micro_utils.coroutines
|
||||
|
||||
import kotlinx.coroutines.CancellationException
|
||||
import kotlinx.coroutines.NonCancellable
|
||||
import kotlinx.coroutines.withContext
|
||||
import kotlin.contracts.ExperimentalContracts
|
||||
import kotlin.contracts.InvocationKind
|
||||
import kotlin.contracts.contract
|
||||
@@ -21,6 +23,7 @@ class SmartRWLocker(private val readPermits: Int = Int.MAX_VALUE, writeIsLocked:
|
||||
val readSemaphore: SmartSemaphore.Immutable = _readSemaphore.immutable()
|
||||
val writeMutex: SmartMutex.Immutable = _writeMutex.immutable()
|
||||
|
||||
|
||||
/**
|
||||
* Do lock in [readSemaphore] inside of [writeMutex] locking
|
||||
*/
|
||||
@@ -32,8 +35,8 @@ class SmartRWLocker(private val readPermits: Int = Int.MAX_VALUE, writeIsLocked:
|
||||
/**
|
||||
* Release one read permit in [readSemaphore]
|
||||
*/
|
||||
suspend fun releaseRead(): Boolean {
|
||||
return _readSemaphore.release()
|
||||
suspend fun releaseRead(): Boolean = withContext(NonCancellable) {
|
||||
_readSemaphore.release()
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -44,7 +47,9 @@ class SmartRWLocker(private val readPermits: Int = Int.MAX_VALUE, writeIsLocked:
|
||||
try {
|
||||
_readSemaphore.acquire(readPermits)
|
||||
} catch (e: CancellationException) {
|
||||
_writeMutex.unlock()
|
||||
withContext(NonCancellable) {
|
||||
_writeMutex.unlock()
|
||||
}
|
||||
throw e
|
||||
}
|
||||
}
|
||||
@@ -52,9 +57,9 @@ class SmartRWLocker(private val readPermits: Int = Int.MAX_VALUE, writeIsLocked:
|
||||
/**
|
||||
* Unlock [writeMutex]
|
||||
*/
|
||||
suspend fun unlockWrite(): Boolean {
|
||||
return _writeMutex.unlock().also {
|
||||
if (it) {
|
||||
suspend fun unlockWrite(): Boolean = withContext(NonCancellable) {
|
||||
_writeMutex.unlock().also { unlocked ->
|
||||
if (unlocked) {
|
||||
_readSemaphore.release(readPermits)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -9,6 +9,7 @@ import kotlin.test.assertEquals
|
||||
import kotlin.test.assertFails
|
||||
import kotlin.test.assertFalse
|
||||
import kotlin.test.assertTrue
|
||||
import kotlin.time.Duration.Companion.days
|
||||
import kotlin.time.Duration.Companion.seconds
|
||||
|
||||
class SmartRWLockerTests {
|
||||
@@ -109,6 +110,59 @@ class SmartRWLockerTests {
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
fun failureOnReadFreeingRead() = runTest {
|
||||
val locker = SmartRWLocker()
|
||||
val job = launch {
|
||||
locker.withReadAcquire {
|
||||
while (isActive) {
|
||||
delay(1.days)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
locker.readSemaphore.permitsStateFlow.first {
|
||||
it == locker.readSemaphore.maxPermits - 1
|
||||
}
|
||||
job.cancelAndJoin()
|
||||
locker.readSemaphore.permitsStateFlow.first {
|
||||
it == locker.readSemaphore.maxPermits
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
fun cancelledReaderReleasesPermitUnderContention() = runTest(timeout = 5.seconds) {
|
||||
val locker = SmartRWLocker(readPermits = 2)
|
||||
val reader = launch(Dispatchers.Unconfined) {
|
||||
locker.withReadAcquire {
|
||||
awaitCancellation()
|
||||
}
|
||||
}
|
||||
assertEquals(1, locker.readSemaphore.freePermits)
|
||||
|
||||
// Observe the second acquisition synchronously while it still holds the
|
||||
// semaphore's internal mutex. Cancelling the unconfined reader makes its
|
||||
// cleanup contend for that mutex before the acquisition can release it.
|
||||
val cancellation = launch(Dispatchers.Unconfined) {
|
||||
locker.readSemaphore.permitsStateFlow.first { it == 0 }
|
||||
reader.cancel()
|
||||
}
|
||||
|
||||
// Keep this acquisition on the normal test dispatcher: making it
|
||||
// unconfined would change the ordering that forces cleanup contention.
|
||||
locker.withReadAcquire {
|
||||
cancellation.join()
|
||||
reader.join()
|
||||
assertTrue(reader.isCancelled)
|
||||
assertEquals(
|
||||
1,
|
||||
locker.readSemaphore.freePermits,
|
||||
"The cancelled reader must release its permit while the other reader still holds one"
|
||||
)
|
||||
}
|
||||
assertEquals(2, locker.readSemaphore.freePermits)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun simpleWithReadAcquireTest() {
|
||||
val locker = SmartRWLocker()
|
||||
|
||||
Reference in New Issue
Block a user