SmartMutex and SmartSemaphore fixes

This commit is contained in:
2026-09-23 15:25:35 +06:00
parent 0e98281fa3
commit 89eb46ed6b
5 changed files with 196 additions and 9 deletions

View File

@@ -5,6 +5,10 @@
* `Coroutines`: * `Coroutines`:
* `SmartRWLocker`: * `SmartRWLocker`:
* Fix of `unlockWrite`, `lockWrite` and `releaseRead` calls to pass correct number of permits * Fix of `unlockWrite`, `lockWrite` and `releaseRead` calls to pass correct number of permits
* `SmartMutex`:
* Fix `unlock` call
* `SmartSemaphore`:
* Fix same issues to avoid cancellation exceptions handling errors and several other problems
## 0.25.8.2 ## 0.25.8.2

View File

@@ -1,5 +1,6 @@
package dev.inmo.micro_utils.coroutines package dev.inmo.micro_utils.coroutines
import kotlinx.coroutines.NonCancellable
import kotlinx.coroutines.currentCoroutineContext import kotlinx.coroutines.currentCoroutineContext
import kotlinx.coroutines.flow.StateFlow import kotlinx.coroutines.flow.StateFlow
import kotlinx.coroutines.flow.asStateFlow import kotlinx.coroutines.flow.asStateFlow
@@ -7,6 +8,7 @@ import kotlinx.coroutines.flow.first
import kotlinx.coroutines.isActive import kotlinx.coroutines.isActive
import kotlinx.coroutines.sync.Mutex import kotlinx.coroutines.sync.Mutex
import kotlinx.coroutines.sync.withLock import kotlinx.coroutines.sync.withLock
import kotlinx.coroutines.withContext
import kotlin.contracts.ExperimentalContracts import kotlin.contracts.ExperimentalContracts
import kotlin.contracts.InvocationKind import kotlin.contracts.InvocationKind
import kotlin.contracts.contract import kotlin.contracts.contract
@@ -92,8 +94,8 @@ sealed interface SmartMutex {
* If [isLocked] == true - will change it to false and return true. If current call will not unlock this * If [isLocked] == true - will change it to false and return true. If current call will not unlock this
* [SmartMutex] - false * [SmartMutex] - false
*/ */
suspend fun unlock(): Boolean { suspend fun unlock(): Boolean = withContext(NonCancellable) {
return if (_lockStateFlow.value) { if (_lockStateFlow.value) {
internalChangesMutex.withLock { internalChangesMutex.withLock {
if (_lockStateFlow.value) { if (_lockStateFlow.value) {
_lockStateFlow.value = false _lockStateFlow.value = false

View File

@@ -1,5 +1,6 @@
package dev.inmo.micro_utils.coroutines package dev.inmo.micro_utils.coroutines
import kotlinx.coroutines.NonCancellable
import kotlinx.coroutines.currentCoroutineContext import kotlinx.coroutines.currentCoroutineContext
import kotlinx.coroutines.flow.StateFlow import kotlinx.coroutines.flow.StateFlow
import kotlinx.coroutines.flow.asStateFlow import kotlinx.coroutines.flow.asStateFlow
@@ -8,6 +9,7 @@ import kotlinx.coroutines.isActive
import kotlinx.coroutines.sync.Mutex import kotlinx.coroutines.sync.Mutex
import kotlinx.coroutines.sync.Semaphore import kotlinx.coroutines.sync.Semaphore
import kotlinx.coroutines.sync.withLock import kotlinx.coroutines.sync.withLock
import kotlinx.coroutines.withContext
import kotlin.contracts.ExperimentalContracts import kotlin.contracts.ExperimentalContracts
import kotlin.contracts.InvocationKind import kotlin.contracts.InvocationKind
import kotlin.contracts.contract import kotlin.contracts.contract
@@ -76,7 +78,9 @@ sealed interface SmartSemaphore {
} }
} while (shouldContinue && currentCoroutineContext().isActive) } while (shouldContinue && currentCoroutineContext().isActive)
} catch (e: Throwable) { } catch (e: Throwable) {
release(acquiredPermits) if (acquiredPermits > 0) {
release(acquiredPermits)
}
throw e throw e
} }
} }
@@ -107,9 +111,9 @@ sealed interface SmartSemaphore {
*/ */
suspend fun tryAcquire(permits: Int = 1): Boolean { suspend fun tryAcquire(permits: Int = 1): Boolean {
val checkedPermits = checkedPermits(permits) val checkedPermits = checkedPermits(permits)
return if (_freePermitsStateFlow.value < checkedPermits) { return if (_freePermitsStateFlow.value >= checkedPermits) {
internalChangesMutex.withLock { internalChangesMutex.withLock {
if (_freePermitsStateFlow.value < checkedPermits) { if (_freePermitsStateFlow.value >= checkedPermits) {
_freePermitsStateFlow.value -= checkedPermits _freePermitsStateFlow.value -= checkedPermits
true true
} else { } else {
@@ -125,12 +129,12 @@ sealed interface SmartSemaphore {
* If [freePermits] == true - will change it to false and return true. If current call will not unlock this * If [freePermits] == true - will change it to false and return true. If current call will not unlock this
* [SmartSemaphore] - false * [SmartSemaphore] - false
*/ */
suspend fun release(permits: Int = 1): Boolean { suspend fun release(permits: Int = 1): Boolean = withContext(NonCancellable) {
val checkedPermits = checkedPermits(permits) val checkedPermits = checkedPermits(permits)
return if (_freePermitsStateFlow.value < this.maxPermits) { if (_freePermitsStateFlow.value < maxPermits) {
internalChangesMutex.withLock { internalChangesMutex.withLock {
if (_freePermitsStateFlow.value < this.maxPermits) { if (_freePermitsStateFlow.value < maxPermits) {
_freePermitsStateFlow.value = minOf(_freePermitsStateFlow.value + checkedPermits, this.maxPermits) _freePermitsStateFlow.value = minOf(_freePermitsStateFlow.value + checkedPermits, maxPermits)
true true
} else { } else {
false false

View File

@@ -0,0 +1,73 @@
import dev.inmo.micro_utils.coroutines.SmartMutex
import dev.inmo.micro_utils.coroutines.withLock
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.awaitCancellation
import kotlinx.coroutines.cancel
import kotlinx.coroutines.cancelAndJoin
import kotlinx.coroutines.currentCoroutineContext
import kotlinx.coroutines.flow.first
import kotlinx.coroutines.launch
import kotlinx.coroutines.test.runTest
import kotlin.test.Test
import kotlin.test.assertFalse
import kotlin.test.assertTrue
import kotlin.time.Duration.Companion.seconds
class SmartMutexTests {
@Test
fun cancelledUnlockCompletesUnderContention() = runTest(timeout = 5.seconds) {
val mutex = SmartMutex.Mutable()
// Delegate this acquisition's release to another coroutine. The
// unconfined collector runs while lock() still holds its internal mutex.
val releaser = launch(Dispatchers.Unconfined) {
mutex.lockStateFlow.first { it }
currentCoroutineContext().cancel()
mutex.unlock()
}
// Keep acquisition on the normal test dispatcher so the collector
// attempts cancelled cleanup before the internal mutex is released.
mutex.lock()
releaser.join()
assertTrue(releaser.isCancelled)
assertFalse(mutex.isLocked, "Cancellation must not prevent the delegated unlock")
assertTrue(mutex.tryLock())
assertTrue(mutex.unlock())
}
@Test
fun cancellingWithLockBodyReleasesMutex() = runTest(timeout = 5.seconds) {
val mutex = SmartMutex.Mutable()
val holder = launch(Dispatchers.Unconfined) {
mutex.withLock {
awaitCancellation()
}
}
assertTrue(mutex.isLocked)
holder.cancelAndJoin()
assertFalse(mutex.isLocked)
}
@Test
fun cancelledWaiterDoesNotEnterOrReleaseHeldMutex() = runTest(timeout = 5.seconds) {
val mutex = SmartMutex.Mutable()
var entered = false
mutex.withLock {
val waiter = launch(Dispatchers.Unconfined) {
mutex.withLock {
entered = true
}
}
waiter.cancelAndJoin()
assertFalse(entered)
assertTrue(mutex.isLocked)
}
assertFalse(mutex.isLocked)
}
}

View File

@@ -0,0 +1,104 @@
import dev.inmo.micro_utils.coroutines.SmartSemaphore
import dev.inmo.micro_utils.coroutines.withAcquire
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.awaitCancellation
import kotlinx.coroutines.cancelAndJoin
import kotlinx.coroutines.flow.first
import kotlinx.coroutines.launch
import kotlinx.coroutines.test.runTest
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertFalse
import kotlin.test.assertTrue
import kotlin.time.Duration.Companion.seconds
class SmartSemaphoreTests {
@Test
fun cancelledHolderReleasesPermitUnderContention() = runTest(timeout = 5.seconds) {
val semaphore = SmartSemaphore.Mutable(permits = 2)
val holder = launch(Dispatchers.Unconfined) {
semaphore.withAcquire {
awaitCancellation()
}
}
assertEquals(1, semaphore.freePermits)
// The synchronous observer cancels the holder while the second
// acquisition still owns the semaphore's internal changes mutex.
val cancellation = launch(Dispatchers.Unconfined) {
semaphore.permitsStateFlow.first { it == 0 }
holder.cancel()
}
semaphore.withAcquire {
cancellation.join()
holder.join()
assertTrue(holder.isCancelled)
assertEquals(1, semaphore.freePermits, "The cancelled holder must return its permit")
}
assertEquals(2, semaphore.freePermits)
}
@Test
fun cancelledAcquireReturnsPartialPermitsUnderContention() = runTest(timeout = 5.seconds) {
// One permit belongs to another holder; the waiter can acquire two
// permits immediately, but must wait for the third.
val semaphore = SmartSemaphore.Mutable(permits = 3, acquiredPermits = 1)
lateinit var waiter: kotlinx.coroutines.Job
val cancellation = launch(Dispatchers.Unconfined) {
semaphore.permitsStateFlow.first { it == 1 }
waiter.cancel()
}
waiter = launch(Dispatchers.Unconfined) {
semaphore.acquire(3)
}
assertEquals(0, semaphore.freePermits)
assertFalse(waiter.isCompleted)
// Publishing this release resumes the observer while the internal
// mutex is held. The cancelled acquire must wait to roll back safely.
semaphore.release()
cancellation.join()
waiter.join()
assertTrue(waiter.isCancelled)
assertEquals(3, semaphore.freePermits, "Cancellation must return both partially acquired permits")
semaphore.withAcquire(3) {
assertEquals(0, semaphore.freePermits)
}
assertEquals(3, semaphore.freePermits)
}
@Test
fun cancelledAcquireWithoutPermitsDoesNotReleaseAnotherHoldersPermit() = runTest(timeout = 5.seconds) {
val semaphore = SmartSemaphore.Mutable(permits = 1, acquiredPermits = 1)
val waiter = launch(Dispatchers.Unconfined) {
semaphore.acquire()
}
assertFalse(waiter.isCompleted)
waiter.cancelAndJoin()
assertEquals(0, semaphore.freePermits, "A cancelled waiter that acquired nothing must release nothing")
assertTrue(semaphore.release())
assertEquals(1, semaphore.freePermits)
}
@Test
fun tryAcquireUsesAvailablePermits() = runTest {
val semaphore = SmartSemaphore.Mutable(permits = 3)
assertTrue(semaphore.tryAcquire(2))
assertEquals(1, semaphore.freePermits)
assertTrue(semaphore.tryAcquire())
assertEquals(0, semaphore.freePermits)
assertTrue(semaphore.release(3))
assertEquals(3, semaphore.freePermits)
}
@Test
fun tryAcquireWithInsufficientPermitsLeavesStateUnchanged() = runTest {
val semaphore = SmartSemaphore.Mutable(permits = 3, acquiredPermits = 2)
assertFalse(semaphore.tryAcquire(2))
assertEquals(1, semaphore.freePermits)
semaphore.acquire()
assertFalse(semaphore.tryAcquire())
assertEquals(0, semaphore.freePermits)
}
}