mirror of
https://github.com/InsanusMokrassar/MicroUtils.git
synced 2026-09-27 02:45:12 +00:00
Compare commits
8 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| c7f603c1bc | |||
| 2633263093 | |||
| c4f15c32ea | |||
| a070f7d2f3 | |||
| 74501079c3 | |||
| 761cfffae1 | |||
| f26fb352e5 | |||
| b71781727a |
27
CHANGELOG.md
27
CHANGELOG.md
@@ -1,8 +1,35 @@
|
||||
# Changelog
|
||||
|
||||
## 0.32.0
|
||||
|
||||
* `Versions`:
|
||||
* `Kotlin`: `2.4.10` -> `2.4.20`
|
||||
* `KSLog`: `1.7.0` -> `2.1.0`
|
||||
* `Compose`: `1.12.0` -> `1.12.1`
|
||||
* `Ktor`: `3.5.2` -> `3.6.0`
|
||||
* `Okio`: `3.18.1` -> `3.18.2`
|
||||
* `KSP`: `2.3.11` -> `2.3.12`
|
||||
* `KotlinPoet`: `2.3.0` -> `2.4.0`
|
||||
* `Gradle Versions`: `0.61.0` -> `0.64.0`
|
||||
* `NMCP`: `1.6.1` -> `1.6.2`
|
||||
* `AndroidX Core KTX`: `1.19.0` -> `1.19.1`
|
||||
* `AndroidX Fragment`: `1.9.0` -> `1.9.1`
|
||||
* `crypto-js`: `4.1.1` -> `4.2.0`
|
||||
|
||||
## 0.31.1
|
||||
|
||||
* `Coroutines`:
|
||||
* `SmartRWLocker`:
|
||||
* 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.31.0
|
||||
|
||||
* `Versions`:
|
||||
* `Kotlin`: `2.3.21` -> `2.4.10`
|
||||
* `SQLite`: `3.53.2.1` -> `3.53.4.0`
|
||||
* `KSP`: `2.3.9` -> `2.3.11`
|
||||
* `Exposed`: `1.3.0` -> `1.5.0`
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
package dev.inmo.micro_utils.coroutines
|
||||
|
||||
import kotlinx.coroutines.NonCancellable
|
||||
import kotlinx.coroutines.currentCoroutineContext
|
||||
import kotlinx.coroutines.flow.StateFlow
|
||||
import kotlinx.coroutines.flow.asStateFlow
|
||||
@@ -7,6 +8,7 @@ import kotlinx.coroutines.flow.first
|
||||
import kotlinx.coroutines.isActive
|
||||
import kotlinx.coroutines.sync.Mutex
|
||||
import kotlinx.coroutines.sync.withLock
|
||||
import kotlinx.coroutines.withContext
|
||||
import kotlin.contracts.ExperimentalContracts
|
||||
import kotlin.contracts.InvocationKind
|
||||
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
|
||||
* [SmartMutex] - false
|
||||
*/
|
||||
suspend fun unlock(): Boolean {
|
||||
return if (_lockStateFlow.value) {
|
||||
suspend fun unlock(): Boolean = withContext(NonCancellable) {
|
||||
if (_lockStateFlow.value) {
|
||||
internalChangesMutex.withLock {
|
||||
if (_lockStateFlow.value) {
|
||||
_lockStateFlow.value = false
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
package dev.inmo.micro_utils.coroutines
|
||||
|
||||
import kotlinx.coroutines.NonCancellable
|
||||
import kotlinx.coroutines.currentCoroutineContext
|
||||
import kotlinx.coroutines.flow.StateFlow
|
||||
import kotlinx.coroutines.flow.asStateFlow
|
||||
@@ -8,6 +9,7 @@ import kotlinx.coroutines.isActive
|
||||
import kotlinx.coroutines.sync.Mutex
|
||||
import kotlinx.coroutines.sync.Semaphore
|
||||
import kotlinx.coroutines.sync.withLock
|
||||
import kotlinx.coroutines.withContext
|
||||
import kotlin.contracts.ExperimentalContracts
|
||||
import kotlin.contracts.InvocationKind
|
||||
import kotlin.contracts.contract
|
||||
@@ -76,7 +78,9 @@ sealed interface SmartSemaphore {
|
||||
}
|
||||
} while (shouldContinue && currentCoroutineContext().isActive)
|
||||
} catch (e: Throwable) {
|
||||
release(acquiredPermits)
|
||||
if (acquiredPermits > 0) {
|
||||
release(acquiredPermits)
|
||||
}
|
||||
throw e
|
||||
}
|
||||
}
|
||||
@@ -107,9 +111,9 @@ sealed interface SmartSemaphore {
|
||||
*/
|
||||
suspend fun tryAcquire(permits: Int = 1): Boolean {
|
||||
val checkedPermits = checkedPermits(permits)
|
||||
return if (_freePermitsStateFlow.value < checkedPermits) {
|
||||
return if (_freePermitsStateFlow.value >= checkedPermits) {
|
||||
internalChangesMutex.withLock {
|
||||
if (_freePermitsStateFlow.value < checkedPermits) {
|
||||
if (_freePermitsStateFlow.value >= checkedPermits) {
|
||||
_freePermitsStateFlow.value -= checkedPermits
|
||||
true
|
||||
} 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
|
||||
* [SmartSemaphore] - false
|
||||
*/
|
||||
suspend fun release(permits: Int = 1): Boolean {
|
||||
suspend fun release(permits: Int = 1): Boolean = withContext(NonCancellable) {
|
||||
val checkedPermits = checkedPermits(permits)
|
||||
return if (_freePermitsStateFlow.value < this.maxPermits) {
|
||||
if (_freePermitsStateFlow.value < maxPermits) {
|
||||
internalChangesMutex.withLock {
|
||||
if (_freePermitsStateFlow.value < this.maxPermits) {
|
||||
_freePermitsStateFlow.value = minOf(_freePermitsStateFlow.value + checkedPermits, this.maxPermits)
|
||||
if (_freePermitsStateFlow.value < maxPermits) {
|
||||
_freePermitsStateFlow.value = minOf(_freePermitsStateFlow.value + checkedPermits, maxPermits)
|
||||
true
|
||||
} else {
|
||||
false
|
||||
|
||||
73
coroutines/src/commonTest/kotlin/SmartMutexTests.kt
Normal file
73
coroutines/src/commonTest/kotlin/SmartMutexTests.kt
Normal 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)
|
||||
}
|
||||
}
|
||||
@@ -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()
|
||||
|
||||
104
coroutines/src/commonTest/kotlin/SmartSemaphoreTests.kt
Normal file
104
coroutines/src/commonTest/kotlin/SmartSemaphoreTests.kt
Normal 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)
|
||||
}
|
||||
}
|
||||
@@ -11,10 +11,10 @@ org.gradle.jvmargs=-Xmx2g -XX:MaxMetaspaceSize=2g
|
||||
|
||||
# JS NPM
|
||||
|
||||
crypto_js_version=4.1.1
|
||||
crypto_js_version=4.2.0
|
||||
|
||||
# Project data
|
||||
|
||||
group=dev.inmo
|
||||
version=0.31.0
|
||||
android_code_version=317
|
||||
version=0.32.0
|
||||
android_code_version=319
|
||||
|
||||
@@ -1,14 +1,14 @@
|
||||
[versions]
|
||||
|
||||
kt = "2.4.10"
|
||||
kt = "2.4.20"
|
||||
kt-serialization = "1.11.0"
|
||||
kt-coroutines = "1.11.0"
|
||||
|
||||
kotlinx-browser = "0.5.0"
|
||||
|
||||
kslog = "1.7.0"
|
||||
kslog = "2.1.0"
|
||||
|
||||
jb-compose = "1.12.0"
|
||||
jb-compose = "1.12.1"
|
||||
jb-compose-material3 = "1.11.0-alpha07"
|
||||
jb-compose-icons = "1.7.8"
|
||||
jb-exposed = "1.5.0"
|
||||
@@ -19,27 +19,27 @@ sqlite = "3.53.4.0"
|
||||
korlibs = "5.4.0"
|
||||
uuid = "0.8.4"
|
||||
|
||||
ktor = "3.5.2"
|
||||
ktor = "3.6.0"
|
||||
|
||||
gh-release = "2.5.2"
|
||||
|
||||
koin = "4.2.2"
|
||||
|
||||
okio = "3.18.1"
|
||||
okio = "3.18.2"
|
||||
|
||||
ksp = "2.3.11"
|
||||
kotlin-poet = "2.3.0"
|
||||
ksp = "2.3.12"
|
||||
kotlin-poet = "2.4.0"
|
||||
|
||||
versions = "0.61.0"
|
||||
nmcp = "1.6.1"
|
||||
versions = "0.64.0"
|
||||
nmcp = "1.6.2"
|
||||
|
||||
android-gradle = "9.1.1"
|
||||
dexcount = "4.0.0"
|
||||
|
||||
android-coreKtx = "1.19.0"
|
||||
android-coreKtx = "1.19.1"
|
||||
android-recyclerView = "1.4.0"
|
||||
android-appCompat = "1.8.0"
|
||||
android-fragment = "1.9.0"
|
||||
android-fragment = "1.9.1"
|
||||
android-espresso = "3.7.0"
|
||||
android-test = "1.3.0"
|
||||
|
||||
|
||||
Reference in New Issue
Block a user