From 0e98281fa3eb0a8433c84af0c351ce4f817efd96 Mon Sep 17 00:00:00 2001 From: InsanusMokrassar Date: Wed, 23 Sep 2026 14:45:36 +0600 Subject: [PATCH] fixes in smart rw locker --- CHANGELOG.md | 4 ++ .../micro_utils/coroutines/SmartRWLocker.kt | 17 +++--- .../commonTest/kotlin/SmartRWLockerTests.kt | 54 +++++++++++++++++++ 3 files changed, 69 insertions(+), 6 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 97084333ba5..9a1b333cb44 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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`: diff --git a/coroutines/src/commonMain/kotlin/dev/inmo/micro_utils/coroutines/SmartRWLocker.kt b/coroutines/src/commonMain/kotlin/dev/inmo/micro_utils/coroutines/SmartRWLocker.kt index aea220812cd..f1f74e7b59f 100644 --- a/coroutines/src/commonMain/kotlin/dev/inmo/micro_utils/coroutines/SmartRWLocker.kt +++ b/coroutines/src/commonMain/kotlin/dev/inmo/micro_utils/coroutines/SmartRWLocker.kt @@ -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) } } diff --git a/coroutines/src/commonTest/kotlin/SmartRWLockerTests.kt b/coroutines/src/commonTest/kotlin/SmartRWLockerTests.kt index 210ce222b15..de4b374a9a3 100644 --- a/coroutines/src/commonTest/kotlin/SmartRWLockerTests.kt +++ b/coroutines/src/commonTest/kotlin/SmartRWLockerTests.kt @@ -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()