diff --git a/fleet/kernel/srcCommonMain/fleet/kernel/rebase/Memoizer.kt b/fleet/kernel/srcCommonMain/fleet/kernel/rebase/Memoizer.kt index 9cf5c3dc47da..a45af0cc80bb 100644 --- a/fleet/kernel/srcCommonMain/fleet/kernel/rebase/Memoizer.kt +++ b/fleet/kernel/srcCommonMain/fleet/kernel/rebase/Memoizer.kt @@ -1,12 +1,12 @@ // Copyright 2000-2024 JetBrains s.r.o. and contributors. Use of this source code is governed by the Apache 2.0 license. package fleet.kernel.rebase -import fleet.multiplatform.shims.ConcurrentHashMap +import fleet.multiplatform.shims.MultiplatformConcurrentHashMap internal class Memoizer { data class WithEpoch(val value: T, val epoch: Int) - private val m = ConcurrentHashMap>() + private val m = MultiplatformConcurrentHashMap>() private var epoch: Int = 0 fun nextEpoch() { diff --git a/fleet/lsp.protocol/srcCommonMain/com/jetbrains/lsp/implementation/withLsp.kt b/fleet/lsp.protocol/srcCommonMain/com/jetbrains/lsp/implementation/withLsp.kt index 7e713e7fdcca..753ad0a1002f 100644 --- a/fleet/lsp.protocol/srcCommonMain/com/jetbrains/lsp/implementation/withLsp.kt +++ b/fleet/lsp.protocol/srcCommonMain/com/jetbrains/lsp/implementation/withLsp.kt @@ -1,7 +1,7 @@ package com.jetbrains.lsp.implementation import com.jetbrains.lsp.protocol.* -import fleet.multiplatform.shims.ConcurrentHashMap +import fleet.multiplatform.shims.MultiplatformConcurrentHashMap import fleet.util.logging.logger import kotlinx.coroutines.* import kotlinx.coroutines.channels.ReceiveChannel @@ -41,7 +41,7 @@ suspend fun withLsp( body: suspend CoroutineScope.(LspClient) -> Unit, ) { coroutineScope { - val outgoingRequests = ConcurrentHashMap() + val outgoingRequests = MultiplatformConcurrentHashMap() val idGen = AtomicInt(0) val lspClient = object : LspClient { override suspend fun request( @@ -104,7 +104,7 @@ suspend fun withLsp( launch(createCoroutineContext(lspClient)) { withSupervisor { supervisor -> - val incomingRequestsJobs = ConcurrentHashMap() + val incomingRequestsJobs = MultiplatformConcurrentHashMap() incoming.consumeEach { jsonMessage -> when { jsonMessage !is JsonObject || jsonMessage["jsonrpc"] != JsonPrimitive("2.0") -> { diff --git a/fleet/multiplatform.shims/srcCommonMain/fleet/multiplatform/shims/ConcurrentHashMap.kt b/fleet/multiplatform.shims/srcCommonMain/fleet/multiplatform/shims/MultiplatformConcurrentHashMap.kt similarity index 63% rename from fleet/multiplatform.shims/srcCommonMain/fleet/multiplatform/shims/ConcurrentHashMap.kt rename to fleet/multiplatform.shims/srcCommonMain/fleet/multiplatform/shims/MultiplatformConcurrentHashMap.kt index e467dd3f86ac..705a76066dc0 100644 --- a/fleet/multiplatform.shims/srcCommonMain/fleet/multiplatform/shims/ConcurrentHashMap.kt +++ b/fleet/multiplatform.shims/srcCommonMain/fleet/multiplatform/shims/MultiplatformConcurrentHashMap.kt @@ -6,11 +6,11 @@ import fleet.util.multiplatform.linkToActual /** * On JVM, this is a wrapper for [java.util.concurrent.ConcurrentHashMap] * - * @see [ConcurrentHashMapJvm] and [ConcurrentHashMapWasmJs] actual implementations + * @see [MultiplatformMultiplatformConcurrentHashMapJvmImpl] and [MultiplatformConcurrentHashMapWasmJsImpl] actual implementations */ -fun ConcurrentHashMap(): ConcurrentHashMap = linkToActual() +fun MultiplatformConcurrentHashMap(): MultiplatformConcurrentHashMap = linkToActual() -interface ConcurrentHashMap: MutableMap { +interface MultiplatformConcurrentHashMap: MutableMap { fun putIfAbsent(key: K, value: V): V? fun computeIfAbsent(key: K, f: (K) -> V): V diff --git a/fleet/multiplatform.shims/srcJvmMain/fleet/multiplatform/shims/ConcurrentHashMap.kt b/fleet/multiplatform.shims/srcJvmMain/fleet/multiplatform/shims/ConcurrentHashMap.kt new file mode 100644 index 000000000000..587c54c40cf8 --- /dev/null +++ b/fleet/multiplatform.shims/srcJvmMain/fleet/multiplatform/shims/ConcurrentHashMap.kt @@ -0,0 +1,58 @@ +// Copyright 2000-2025 JetBrains s.r.o. and contributors. Use of this source code is governed by the Apache 2.0 license. +package fleet.multiplatform.shims + +import java.util.concurrent.ConcurrentHashMap as JavaConcurrentHashMap + +/** + * @deprecated Use [java.util.concurrent.ConcurrentHashMap] + * + * Factory method is not supposed to be used outside the KMP context. + * + * Moreover, using it may lead to quite tricky bugs while using extension methods on MutableMap, + * that are not ready for concurrent execution, instead of the one defined on ConcurrentHashMap (e.g. `getOrPut function) + */ +@Deprecated("Use ConcurrentHashMap from java.util.concurrent", ReplaceWith("ConcurrentHashMap()", "java.util.concurrent.ConcurrentHashMap")) +fun ConcurrentHashMap(): ConcurrentHashMap = ConcurrentHashMapImpl(JavaConcurrentHashMap()) + +/** + * @deprecated Use [java.util.concurrent.ConcurrentHashMap] + * + * Interface is not supposed to be used outside the KMP context. + * + * Moreover, using it may lead to quite tricky bugs while using extension methods on MutableMap, + * that are not ready for concurrent execution, instead of the one defined on ConcurrentHashMap (e.g. `getOrPut function) + */ +@Deprecated("Use ConcurrentHashMap from java.util.concurrent", ReplaceWith("ConcurrentHashMap", "java.util.concurrent.ConcurrentHashMap")) +interface ConcurrentHashMap : MutableMap { + fun putIfAbsent(key: K, value: V): V? + + fun computeIfAbsent(key: K, f: (K) -> V): V + + fun computeIfPresent(key: K, f: (K, V) -> V): V? + + fun compute(key: K, f: (K, V?) -> V?): V? + + fun remove(key: K, value: V): Boolean +} + +private class ConcurrentHashMapImpl(private val hashMap: JavaConcurrentHashMap) : MutableMap by hashMap, ConcurrentHashMap { + override fun remove(key: K, value: V): Boolean { + return hashMap.remove(key, value) + } + + override fun putIfAbsent(key: K, value: V): V? { + return hashMap.putIfAbsent(key, value) + } + + override fun computeIfAbsent(key: K, f: (K) -> V): V { + return hashMap.computeIfAbsent(key, f) + } + + override fun computeIfPresent(key: K, f: (K, V) -> V): V? { + return hashMap.computeIfPresent(key, f) + } + + override fun compute(key: K, f: (K, V?) -> V?): V? { + return hashMap.compute(key, f) + } +} \ No newline at end of file diff --git a/fleet/multiplatform.shims/srcJvmMain/fleet/multiplatform/shims/ConcurrentHashMap.jvm.kt b/fleet/multiplatform.shims/srcJvmMain/fleet/multiplatform/shims/MultiplatformConcurrentHashMap.jvm.kt similarity index 68% rename from fleet/multiplatform.shims/srcJvmMain/fleet/multiplatform/shims/ConcurrentHashMap.jvm.kt rename to fleet/multiplatform.shims/srcJvmMain/fleet/multiplatform/shims/MultiplatformConcurrentHashMap.jvm.kt index 7700a07500fe..2c3cbaadbad6 100644 --- a/fleet/multiplatform.shims/srcJvmMain/fleet/multiplatform/shims/ConcurrentHashMap.jvm.kt +++ b/fleet/multiplatform.shims/srcJvmMain/fleet/multiplatform/shims/MultiplatformConcurrentHashMap.jvm.kt @@ -6,9 +6,9 @@ import fleet.util.multiplatform.Actual import java.util.concurrent.ConcurrentHashMap as JavaConcurrentHashMap @Actual -internal fun ConcurrentHashMapJvm(): ConcurrentHashMap = MultiplatformConcurrentHashMap(JavaConcurrentHashMap()) +internal fun MultiplatformConcurrentHashMapJvm(): MultiplatformConcurrentHashMap = MultiplatformMultiplatformConcurrentHashMapJvmImpl(JavaConcurrentHashMap()) -private class MultiplatformConcurrentHashMap(val hashMap: JavaConcurrentHashMap) : MutableMap by hashMap, ConcurrentHashMap { +private class MultiplatformMultiplatformConcurrentHashMapJvmImpl(val hashMap: JavaConcurrentHashMap) : MutableMap by hashMap, MultiplatformConcurrentHashMap { override fun remove(key: K, value: V): Boolean { return hashMap.remove(key, value) } diff --git a/fleet/multiplatform.shims/srcWasmJsMain/fleet/multiplatform/shims/ConcurrentHashMap.wasm.kt b/fleet/multiplatform.shims/srcWasmJsMain/fleet/multiplatform/shims/MultiplatformConcurrentHashMap.wasm.kt similarity index 80% rename from fleet/multiplatform.shims/srcWasmJsMain/fleet/multiplatform/shims/ConcurrentHashMap.wasm.kt rename to fleet/multiplatform.shims/srcWasmJsMain/fleet/multiplatform/shims/MultiplatformConcurrentHashMap.wasm.kt index 793fffd53508..0e9b24c8939b 100644 --- a/fleet/multiplatform.shims/srcWasmJsMain/fleet/multiplatform/shims/ConcurrentHashMap.wasm.kt +++ b/fleet/multiplatform.shims/srcWasmJsMain/fleet/multiplatform/shims/MultiplatformConcurrentHashMap.wasm.kt @@ -5,9 +5,9 @@ package fleet.multiplatform.shims import fleet.util.multiplatform.Actual @Actual -fun ConcurrentHashMapWasmJs(): ConcurrentHashMap = ConcurrentHashMapWasm(mutableMapOf()) +fun MultiplatformConcurrentHashMapWasmJs(): MultiplatformConcurrentHashMap = MultiplatformConcurrentHashMapWasmJsImpl(mutableMapOf()) -internal fun ConcurrentHashMapWasm(base: MutableMap): ConcurrentHashMap = object : MutableMap by base, ConcurrentHashMap { +internal fun MultiplatformConcurrentHashMapWasmJsImpl(base: MutableMap): MultiplatformConcurrentHashMap = object : MutableMap by base, MultiplatformConcurrentHashMap { override fun putIfAbsent(key: K, value: V): V? { return if (!containsKey(key)) { put(key, value) diff --git a/fleet/rhizomedb/srcCommonMain/com/jetbrains/rhizomedb/impl/EidGen.kt b/fleet/rhizomedb/srcCommonMain/com/jetbrains/rhizomedb/impl/EidGen.kt index 496d835fdcc6..7b629e2150e1 100644 --- a/fleet/rhizomedb/srcCommonMain/com/jetbrains/rhizomedb/impl/EidGen.kt +++ b/fleet/rhizomedb/srcCommonMain/com/jetbrains/rhizomedb/impl/EidGen.kt @@ -5,7 +5,7 @@ import com.jetbrains.rhizomedb.EID import com.jetbrains.rhizomedb.MAX_PART import com.jetbrains.rhizomedb.Part import com.jetbrains.rhizomedb.withPart -import fleet.multiplatform.shims.ConcurrentHashMap +import fleet.multiplatform.shims.MultiplatformConcurrentHashMap import kotlin.concurrent.atomics.AtomicInt import kotlin.concurrent.atomics.incrementAndFetch @@ -14,7 +14,7 @@ sealed interface EidGen { private val eidGens: Array = Array(MAX_PART + 1) { AtomicInt(0) }.also { it[0].store(17) } - private val eidMemo: ConcurrentHashMap = ConcurrentHashMap() + private val eidMemo: MultiplatformConcurrentHashMap = MultiplatformConcurrentHashMap() override fun freshEID(part: Part): EID = withPart(eidGens[part].incrementAndFetch(), part) diff --git a/fleet/rpc.server/srcCommonMain/fleet/rpc/server/RpcExecutor.kt b/fleet/rpc.server/srcCommonMain/fleet/rpc/server/RpcExecutor.kt index 6d2d5c7bc5ab..4fff0e76e59c 100644 --- a/fleet/rpc.server/srcCommonMain/fleet/rpc/server/RpcExecutor.kt +++ b/fleet/rpc.server/srcCommonMain/fleet/rpc/server/RpcExecutor.kt @@ -18,7 +18,7 @@ import kotlinx.coroutines.channels.Channel import kotlinx.coroutines.channels.SendChannel import kotlinx.coroutines.channels.consumeEach import kotlinx.serialization.json.Json -import fleet.multiplatform.shims.ConcurrentHashMap +import fleet.multiplatform.shims.MultiplatformConcurrentHashMap import fleet.multiplatform.shims.MultiplatformConcurrentHashSet import fleet.util.async.withSupervisor import kotlinx.serialization.builtins.serializer @@ -33,10 +33,10 @@ class RpcExecutor private constructor( private val rpcCallDispatcher: CoroutineDispatcher?, ) { - private val remoteObjects = ConcurrentHashMap() - private val resources = ConcurrentHashMap() - private val children: ConcurrentHashMap> = ConcurrentHashMap() - private val parents: ConcurrentHashMap = ConcurrentHashMap() + private val remoteObjects = MultiplatformConcurrentHashMap() + private val resources = MultiplatformConcurrentHashMap() + private val children: MultiplatformConcurrentHashMap> = MultiplatformConcurrentHashMap() + private val parents: MultiplatformConcurrentHashMap = MultiplatformConcurrentHashMap() companion object { internal val logger = KLoggers.logger(RpcExecutor::class) @@ -100,10 +100,10 @@ class RpcExecutor private constructor( } } - private val requestJobs = ConcurrentHashMap() - private val routeRequests = ConcurrentHashMap>() - private val channels = ConcurrentHashMap() - private val routeChannels = ConcurrentHashMap>() + private val requestJobs = MultiplatformConcurrentHashMap() + private val routeRequests = MultiplatformConcurrentHashMap>() + private val channels = MultiplatformConcurrentHashMap() + private val routeChannels = MultiplatformConcurrentHashMap>() private suspend fun send(message: TransportMessage) { sendSuspend(::sendAsync, message) diff --git a/fleet/rpc.server/srcCommonMain/fleet/rpc/server/ServerRequestDispatcher.kt b/fleet/rpc.server/srcCommonMain/fleet/rpc/server/ServerRequestDispatcher.kt index 10132d852b1f..a7b64cdb025a 100644 --- a/fleet/rpc.server/srcCommonMain/fleet/rpc/server/ServerRequestDispatcher.kt +++ b/fleet/rpc.server/srcCommonMain/fleet/rpc/server/ServerRequestDispatcher.kt @@ -12,7 +12,7 @@ import kotlinx.coroutines.channels.consume import kotlinx.coroutines.channels.consumeEach import kotlinx.coroutines.flow.* import kotlinx.coroutines.selects.select -import fleet.multiplatform.shims.ConcurrentHashMap +import fleet.multiplatform.shims.MultiplatformConcurrentHashMap import fleet.rpc.EndpointKind class ServerRequestDispatcher(private val connectionListener: ConnectionListener?) : RequestDispatcher { @@ -22,7 +22,7 @@ class ServerRequestDispatcher(private val connectionListener: ConnectionListener } private val bannedEndpoints = MutableStateFlow>(emptySet()) - private val connections = ConcurrentHashMap>() + private val connections = MultiplatformConcurrentHashMap>() fun ban(route: UID) { bannedEndpoints.update { it + route } diff --git a/fleet/rpc/srcCommonMain/fleet/rpc/client/RpcClient.kt b/fleet/rpc/srcCommonMain/fleet/rpc/client/RpcClient.kt index 934f52287b9b..a19ac8bfae0a 100644 --- a/fleet/rpc/srcCommonMain/fleet/rpc/client/RpcClient.kt +++ b/fleet/rpc/srcCommonMain/fleet/rpc/client/RpcClient.kt @@ -1,7 +1,7 @@ // Copyright 2000-2024 JetBrains s.r.o. and contributors. Use of this source code is governed by the Apache 2.0 license. package fleet.rpc.client -import fleet.multiplatform.shims.ConcurrentHashMap +import fleet.multiplatform.shims.MultiplatformConcurrentHashMap import fleet.multiplatform.shims.newSingleThreadCoroutineDispatcher import fleet.rpc.RemoteApiDescriptor import fleet.rpc.RemoteKind @@ -81,13 +81,13 @@ private class RpcClient( element.second?.invoke(RpcClientDisconnectedException("Request channel closed", null)) } - private val grayList = ConcurrentHashMap>() + private val grayList = MultiplatformConcurrentHashMap>() - private val outgoingRpc = ConcurrentHashMap() - private val completedRpc = ConcurrentHashMap() - private val streams = ConcurrentHashMap() - private val remoteResources = ConcurrentHashMap>>() - private val resourceParents = ConcurrentHashMap() + private val outgoingRpc = MultiplatformConcurrentHashMap() + private val completedRpc = MultiplatformConcurrentHashMap() + private val streams = MultiplatformConcurrentHashMap() + private val remoteResources = MultiplatformConcurrentHashMap>>() + private val resourceParents = MultiplatformConcurrentHashMap() private val remoteObjectFactory = this.asHandlerFactory().tracing() diff --git a/fleet/rpc/srcCommonMain/fleet/rpc/client/proxy/ProxyCache.kt b/fleet/rpc/srcCommonMain/fleet/rpc/client/proxy/ProxyCache.kt index 03ad9bb3107c..cf12630f2e2b 100644 --- a/fleet/rpc/srcCommonMain/fleet/rpc/client/proxy/ProxyCache.kt +++ b/fleet/rpc/srcCommonMain/fleet/rpc/client/proxy/ProxyCache.kt @@ -5,7 +5,7 @@ import fleet.rpc.RemoteApi import fleet.rpc.RemoteApiDescriptor import fleet.rpc.core.InstanceId import fleet.util.UID -import fleet.multiplatform.shims.ConcurrentHashMap +import fleet.multiplatform.shims.MultiplatformConcurrentHashMap interface ProxyCache { fun proxy(remoteApiDescriptor: RemoteApiDescriptor<*>, key: K, proxy: () -> T): T @@ -16,7 +16,7 @@ fun proxyCache(): ProxyCache { data class ProxyCacheValue(val proxy: Any, val remoteApiDescriptor: RemoteApiDescriptor<*>) - val cache = ConcurrentHashMap() + val cache = MultiplatformConcurrentHashMap() return object : ProxyCache { @Suppress("UNCHECKED_CAST") override fun proxy(remoteApiDescriptor: RemoteApiDescriptor<*>, key: K, proxy: () -> T): T = diff --git a/fleet/util/network/srcJvmMain/fleet/net/TcpProxy.kt b/fleet/util/network/srcJvmMain/fleet/net/TcpProxy.kt index d97cfe9ed2e1..3c87f7451ab1 100644 --- a/fleet/util/network/srcJvmMain/fleet/net/TcpProxy.kt +++ b/fleet/util/network/srcJvmMain/fleet/net/TcpProxy.kt @@ -14,7 +14,7 @@ import kotlinx.coroutines.channels.Channel import kotlinx.coroutines.channels.ReceiveChannel import kotlinx.coroutines.channels.SendChannel import kotlinx.coroutines.channels.consumeEach -import fleet.multiplatform.shims.ConcurrentHashMap +import fleet.multiplatform.shims.MultiplatformConcurrentHashMap import fleet.multiplatform.shims.multiplatformIO private object TcpProxy { @@ -37,7 +37,7 @@ private fun ServerSocket.port(): Int? { } } -private val selectorManagers = ConcurrentHashMap() +private val selectorManagers = MultiplatformConcurrentHashMap() @OptIn(InternalCoroutinesApi::class) // unfortunately necessary to unblock selector manager on cancellation private fun getSelectorManager(parentJob: Job): ActorSelectorManager =