From dbdff0785d982636af2ce83ee794d32069f0375b Mon Sep 17 00:00:00 2001 From: Alexander Zolotov Date: Tue, 7 Oct 2025 12:31:01 +0200 Subject: [PATCH] [fleet, util] replace atomic shims with stdlib functions GitOrigin-RevId: aec14ee003b0b608f801feff7eb1d66e2abfa79d --- .../com/jetbrains/rhizomedb/QueryCache.kt | 6 ++--- .../srcCommonMain/fleet/rpc/core/RpcStream.kt | 12 +++++----- .../fleet/util/AtomicExtensions.kt | 9 -------- .../srcCommonMain/fleet/util/async/Handle.kt | 10 ++++++--- .../util/openmap/MutableBoundedOpenMapImpl.kt | 10 ++++----- .../fleet/util/AtomicExtensions.jvm.kt | 16 -------------- .../fleet/util/AtomicExtensions.wasm.kt | 22 ------------------- 7 files changed, 21 insertions(+), 64 deletions(-) delete mode 100644 fleet/util/core/srcCommonMain/fleet/util/AtomicExtensions.kt delete mode 100644 fleet/util/core/srcJvmMain/fleet/util/AtomicExtensions.jvm.kt delete mode 100644 fleet/util/core/srcWasmJsMain/fleet/util/AtomicExtensions.wasm.kt diff --git a/fleet/rhizomedb/srcCommonMain/com/jetbrains/rhizomedb/QueryCache.kt b/fleet/rhizomedb/srcCommonMain/com/jetbrains/rhizomedb/QueryCache.kt index 76fb950d4b13..c3a4b52f0701 100644 --- a/fleet/rhizomedb/srcCommonMain/com/jetbrains/rhizomedb/QueryCache.kt +++ b/fleet/rhizomedb/srcCommonMain/com/jetbrains/rhizomedb/QueryCache.kt @@ -4,12 +4,12 @@ package com.jetbrains.rhizomedb import fleet.fastutil.longs.LongArrayList import fleet.fastutil.longs.toArray import fleet.util.computeShim -import fleet.util.updateAndGet import kotlinx.collections.immutable.PersistentMap import kotlinx.collections.immutable.PersistentSet import kotlinx.collections.immutable.persistentHashMapOf import kotlinx.collections.immutable.persistentHashSetOf import kotlin.concurrent.atomics.AtomicReference +import kotlin.concurrent.atomics.updateAndFetch class QueryCache private constructor(private val cache: AtomicReference) { @@ -66,14 +66,14 @@ class QueryCache private constructor(private val cache: AtomicReference performQuery(query: CachedQuery, compute: () -> CachedQueryResult): CachedQueryResult = @Suppress("UNCHECKED_CAST") (cache.load().find(query) as CachedQueryResult?) ?: run { val res = compute() - cache.updateAndGet { index -> + cache.updateAndFetch { index -> val existing = index.find(query) when { existing == null -> index.insert(query, res) diff --git a/fleet/rpc/srcCommonMain/fleet/rpc/core/RpcStream.kt b/fleet/rpc/srcCommonMain/fleet/rpc/core/RpcStream.kt index 0549594266fc..f302e3f0c697 100644 --- a/fleet/rpc/srcCommonMain/fleet/rpc/core/RpcStream.kt +++ b/fleet/rpc/srcCommonMain/fleet/rpc/core/RpcStream.kt @@ -3,9 +3,7 @@ package fleet.rpc.core import fleet.util.UID import fleet.util.async.coroutineNameAppended -import fleet.util.getAndUpdate import fleet.util.logging.logger -import fleet.util.updateAndGet import kotlinx.coroutines.* import kotlinx.coroutines.channels.Channel import kotlinx.coroutines.channels.ReceiveChannel @@ -13,6 +11,8 @@ import kotlinx.coroutines.channels.SendChannel import kotlinx.coroutines.channels.consumeEach import kotlinx.serialization.KSerializer import kotlin.concurrent.atomics.AtomicReference +import kotlin.concurrent.atomics.fetchAndUpdate +import kotlin.concurrent.atomics.updateAndFetch import kotlin.coroutines.CoroutineContext import kotlin.coroutines.resumeWithException import kotlin.math.max @@ -54,12 +54,12 @@ class Budget(initial: Int) { private data class State(val cancellation: CancellationException?, val budget: Int, val continuation: CancellableContinuation?) private fun withdraw(): Boolean { - return state.getAndUpdate { it.copy(budget = max(0, it.budget - 1)) }.budget > 0 + return state.fetchAndUpdate { it.copy(budget = max(0, it.budget - 1)) }.budget > 0 } private suspend fun await() { suspendCancellableCoroutine { continuation -> - val r = state.updateAndGet { + val r = state.updateAndFetch { require(it.continuation == null) { "Budget is not intended to use by several producers" } when { it.budget > 0 -> it @@ -82,14 +82,14 @@ class Budget(initial: Int) { fun refill(quantity: Int) { require(quantity > 0) - val r = state.getAndUpdate { + val r = state.fetchAndUpdate { it.copy(budget = it.budget + quantity, continuation = null) } r.continuation?.resumeWith(Result.success(Unit)) } internal fun cancel(cause: CancellationException) { - val was = state.getAndUpdate { + val was = state.fetchAndUpdate { if (it.cancellation == null) { it.copy(cancellation = cause, continuation = null) } diff --git a/fleet/util/core/srcCommonMain/fleet/util/AtomicExtensions.kt b/fleet/util/core/srcCommonMain/fleet/util/AtomicExtensions.kt deleted file mode 100644 index de69794fe491..000000000000 --- a/fleet/util/core/srcCommonMain/fleet/util/AtomicExtensions.kt +++ /dev/null @@ -1,9 +0,0 @@ -// Copyright 2000-2024 JetBrains s.r.o. and contributors. Use of this source code is governed by the Apache 2.0 license. -package fleet.util - -import fleet.util.multiplatform.linkToActual -import kotlin.concurrent.atomics.AtomicReference - -fun AtomicReference.updateAndGet(f: (T) -> T): T = linkToActual() - -fun AtomicReference.getAndUpdate(f: (T) -> T): T = linkToActual() \ No newline at end of file diff --git a/fleet/util/core/srcCommonMain/fleet/util/async/Handle.kt b/fleet/util/core/srcCommonMain/fleet/util/async/Handle.kt index 80584e9bf465..952a6d2eab21 100644 --- a/fleet/util/core/srcCommonMain/fleet/util/async/Handle.kt +++ b/fleet/util/core/srcCommonMain/fleet/util/async/Handle.kt @@ -1,10 +1,10 @@ // Copyright 2000-2024 JetBrains s.r.o. and contributors. Use of this source code is governed by the Apache 2.0 license. package fleet.util.async -import fleet.util.updateAndGet import kotlinx.collections.immutable.persistentSetOf import kotlinx.coroutines.* import kotlin.concurrent.atomics.AtomicReference +import kotlin.concurrent.atomics.updateAndFetch import kotlin.coroutines.CoroutineContext data class Handle(val value: Deferred>, @@ -86,8 +86,12 @@ private suspend fun handleScopeImpl(outerScope: CoroutineScope, body: suspend override val coroutineContext: CoroutineContext get() = context override fun handle(launcher: Launcher): Handle { val handle = outerScope.handle(launcher) - handles.updateAndGet { hs -> hs.add(handle) } - handle.job.invokeOnCompletion { handles.updateAndGet { hs -> hs.remove(handle) } } + handles.updateAndFetch { hs -> hs.add(handle) } + handle.job.invokeOnCompletion { + handles.updateAndFetch { hs -> + hs.remove(handle) + } + } return handle } }.body() diff --git a/fleet/util/core/srcCommonMain/fleet/util/openmap/MutableBoundedOpenMapImpl.kt b/fleet/util/core/srcCommonMain/fleet/util/openmap/MutableBoundedOpenMapImpl.kt index 73ad07024862..784161f035e2 100644 --- a/fleet/util/core/srcCommonMain/fleet/util/openmap/MutableBoundedOpenMapImpl.kt +++ b/fleet/util/core/srcCommonMain/fleet/util/openmap/MutableBoundedOpenMapImpl.kt @@ -1,9 +1,9 @@ // Copyright 2000-2024 JetBrains s.r.o. and contributors. Use of this source code is governed by the Apache 2.0 license. package fleet.util.openmap -import fleet.util.updateAndGet import kotlinx.collections.immutable.PersistentMap import kotlin.concurrent.atomics.AtomicReference +import kotlin.concurrent.atomics.updateAndFetch internal class MutableBoundedOpenMapImpl(val map: AtomicReference, V>>) : MutableBoundedOpenMap { override fun get(k: Key): T? { @@ -13,11 +13,11 @@ internal class MutableBoundedOpenMapImpl(val map: AtomicReferen override fun isEmpty() = map.load().size == 0 override fun set(k: Key, v: T) { - map.updateAndGet { map -> map.put(k, v) } + map.updateAndFetch { map -> map.put(k, v) } } override fun remove(k: Key) { - map.updateAndGet { map -> map.remove(k) } + map.updateAndFetch { map -> map.remove(k) } } override fun assoc(k: Key, v: T): BoundedOpenMap { @@ -33,14 +33,14 @@ internal class MutableBoundedOpenMapImpl(val map: AtomicReferen } override fun update(key: Key, f: (T?) -> T): T { - return map.updateAndGet { map -> + return map.updateAndFetch { map -> val v1 = f(map.get(key) as T?) map.put(key, v1) }.get(key) as T } override fun getOrInit(key: Key, init: () -> T): T { - return map.updateAndGet { map -> + return map.updateAndFetch { map -> val v = map.get(key) if (v == null) { map.put(key, init()) diff --git a/fleet/util/core/srcJvmMain/fleet/util/AtomicExtensions.jvm.kt b/fleet/util/core/srcJvmMain/fleet/util/AtomicExtensions.jvm.kt deleted file mode 100644 index b67dc7f86bf4..000000000000 --- a/fleet/util/core/srcJvmMain/fleet/util/AtomicExtensions.jvm.kt +++ /dev/null @@ -1,16 +0,0 @@ -// Copyright 2000-2025 JetBrains s.r.o. and contributors. Use of this source code is governed by the Apache 2.0 license. -package fleet.util - -import fleet.util.multiplatform.Actual -import kotlin.concurrent.atomics.AtomicReference -import kotlin.concurrent.atomics.asJavaAtomic - -@Actual -fun AtomicReference.updateAndGetJvm(f: (T) -> T): T { - return asJavaAtomic().updateAndGet(f) -} - -@Actual -fun AtomicReference.getAndUpdateJvm(f: (T) -> T): T { - return asJavaAtomic().getAndUpdate(f) -} \ No newline at end of file diff --git a/fleet/util/core/srcWasmJsMain/fleet/util/AtomicExtensions.wasm.kt b/fleet/util/core/srcWasmJsMain/fleet/util/AtomicExtensions.wasm.kt deleted file mode 100644 index 3daf9083ca7a..000000000000 --- a/fleet/util/core/srcWasmJsMain/fleet/util/AtomicExtensions.wasm.kt +++ /dev/null @@ -1,22 +0,0 @@ -// Copyright 2000-2024 JetBrains s.r.o. and contributors. Use of this source code is governed by the Apache 2.0 license. -@file:OptIn(ExperimentalAtomicApi::class) - -package fleet.util - -import fleet.util.multiplatform.Actual -import kotlin.concurrent.atomics.AtomicReference -import kotlin.concurrent.atomics.ExperimentalAtomicApi - -@Actual -fun AtomicReference.updateAndGetWasmJs(f: (T) -> T): T { - val new = f(load()) - store(new) - return new -} - -@Actual -fun AtomicReference.getAndUpdateWasmJs(f: (T) -> T): T { - val old = load() - store(f(old)) - return old -} \ No newline at end of file