diff --git a/fleet/kernel/srcCommonMain/fleet/kernel/rete/impl/Distinct.kt b/fleet/kernel/srcCommonMain/fleet/kernel/rete/impl/Distinct.kt index e55f53a20046..7788d45814d0 100644 --- a/fleet/kernel/srcCommonMain/fleet/kernel/rete/impl/Distinct.kt +++ b/fleet/kernel/srcCommonMain/fleet/kernel/rete/impl/Distinct.kt @@ -2,19 +2,19 @@ package fleet.kernel.rete.impl import fleet.kernel.rete.* -import fleet.multiplatform.shims.ConcurrentHashSet +import fleet.multiplatform.shims.MultiplatformConcurrentHashSet internal fun SubscriptionScope.distinct(producer: Producer): Producer = run { val broadcast = Broadcaster() - val memory = HashMap>>() + val memory = HashMap>>() producer.collect { token -> val v = token.match.value when (token.added) { true -> { when (val matches = memory[v]) { null -> { - val ms = ConcurrentHashSet>() + val ms = MultiplatformConcurrentHashSet>() memory[v] = ms ms.add(token.match) broadcast(Token(true, distinctMatchValue(v, ms))) @@ -29,7 +29,7 @@ internal fun SubscriptionScope.distinct(producer: Producer): Produc if (matches.remove(token.match)) { if (matches.isEmpty()) { memory.remove(v) - broadcast(Token(false, distinctMatchValue(v, emptySet()))) + broadcast(Token(false, distinctMatchValue(v, MultiplatformConcurrentHashSet.empty()))) } } } @@ -44,7 +44,7 @@ internal fun SubscriptionScope.distinct(producer: Producer): Produc } } -private fun distinctMatchValue(v: T, ms: Set>): Match = +private fun distinctMatchValue(v: T, ms: MultiplatformConcurrentHashSet>): Match = Match.validatable(v) { val anyValid = ms.any { match -> match.validate() == ValidationResultEnum.Valid diff --git a/fleet/multiplatform.shims/srcCommonMain/fleet/multiplatform/shims/ConcurrentHashSet.kt b/fleet/multiplatform.shims/srcCommonMain/fleet/multiplatform/shims/ConcurrentHashSet.kt deleted file mode 100644 index 335895299f6b..000000000000 --- a/fleet/multiplatform.shims/srcCommonMain/fleet/multiplatform/shims/ConcurrentHashSet.kt +++ /dev/null @@ -1,7 +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.multiplatform.shims - - -import fleet.util.multiplatform.linkToActual - -fun ConcurrentHashSet(): MutableSet = linkToActual() diff --git a/fleet/multiplatform.shims/srcCommonMain/fleet/multiplatform/shims/EmptyIterator.kt b/fleet/multiplatform.shims/srcCommonMain/fleet/multiplatform/shims/EmptyIterator.kt new file mode 100644 index 000000000000..97ac095df553 --- /dev/null +++ b/fleet/multiplatform.shims/srcCommonMain/fleet/multiplatform/shims/EmptyIterator.kt @@ -0,0 +1,14 @@ +// 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 + +internal object EmptyIterator : MutableListIterator { + override fun hasNext(): Boolean = false + override fun remove() = throw UnsupportedOperationException() + override fun set(element: Nothing) = throw UnsupportedOperationException() + override fun add(element: Nothing) = throw UnsupportedOperationException() + override fun hasPrevious(): Boolean = false + override fun nextIndex(): Int = 0 + override fun previousIndex(): Int = -1 + override fun next(): Nothing = throw NoSuchElementException() + override fun previous(): Nothing = throw NoSuchElementException() +} \ No newline at end of file diff --git a/fleet/multiplatform.shims/srcCommonMain/fleet/multiplatform/shims/MultiplatformConcurrentHashSet.kt b/fleet/multiplatform.shims/srcCommonMain/fleet/multiplatform/shims/MultiplatformConcurrentHashSet.kt new file mode 100644 index 000000000000..becd6ea02cc3 --- /dev/null +++ b/fleet/multiplatform.shims/srcCommonMain/fleet/multiplatform/shims/MultiplatformConcurrentHashSet.kt @@ -0,0 +1,35 @@ +// 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 fleet.util.multiplatform.linkToActual + +fun MultiplatformConcurrentHashSet(): MultiplatformConcurrentHashSet = linkToActual() + +interface MultiplatformConcurrentHashSet : Iterable { + companion object { + fun empty(): MultiplatformConcurrentHashSet { + return EmptySet as MultiplatformConcurrentHashSet + } + } + + val size: Int + fun add(element: @UnsafeVariance T): Boolean + fun addAll(elements: Collection<@UnsafeVariance T>): Boolean + fun remove(element: @UnsafeVariance T): Boolean + operator fun contains(element: @UnsafeVariance T): Boolean + fun isEmpty(): Boolean + fun clear() +} + +private object EmptySet : MultiplatformConcurrentHashSet { + override val size: Int get() = 0 + override fun add(element: Nothing): Boolean = throw UnsupportedOperationException() + override fun addAll(elements: Collection): Boolean = throw UnsupportedOperationException() + override fun remove(element: Nothing): Boolean = false + override fun contains(element: Nothing): Boolean = false + override fun isEmpty(): Boolean = true + override fun clear() = throw UnsupportedOperationException() + override fun iterator(): Iterator = EmptyIterator + override fun equals(other: Any?): Boolean = other is MultiplatformConcurrentHashSet<*> && other.isEmpty() + override fun hashCode(): Int = 0 +} diff --git a/fleet/multiplatform.shims/srcJvmMain/fleet/multiplatform/shims/ConcurrentHashSet.jvm.kt b/fleet/multiplatform.shims/srcJvmMain/fleet/multiplatform/shims/ConcurrentHashSet.jvm.kt deleted file mode 100644 index 57d8a29d9ffc..000000000000 --- a/fleet/multiplatform.shims/srcJvmMain/fleet/multiplatform/shims/ConcurrentHashSet.jvm.kt +++ /dev/null @@ -1,10 +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.multiplatform.shims - - -import fleet.util.multiplatform.Actual -import java.util.concurrent.ConcurrentHashMap as JavaConcurrentHashMap - - -@Actual -internal fun ConcurrentHashSetJvm(): MutableSet = JavaConcurrentHashMap.newKeySet() \ No newline at end of file diff --git a/fleet/multiplatform.shims/srcJvmMain/fleet/multiplatform/shims/ConcurrentHashSet.kt b/fleet/multiplatform.shims/srcJvmMain/fleet/multiplatform/shims/ConcurrentHashSet.kt new file mode 100644 index 000000000000..926b01ca1248 --- /dev/null +++ b/fleet/multiplatform.shims/srcJvmMain/fleet/multiplatform/shims/ConcurrentHashSet.kt @@ -0,0 +1,8 @@ +// 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] instead", replaceWith = ReplaceWith("ConcurrentHashMap.newKeySet()", + "java.util.concurrent.ConcurrentHashMap")) +internal fun ConcurrentHashSet(): MutableSet = JavaConcurrentHashMap.newKeySet() diff --git a/fleet/multiplatform.shims/srcJvmMain/fleet/multiplatform/shims/MultiplatformConcurrentHashSet.jvm.kt b/fleet/multiplatform.shims/srcJvmMain/fleet/multiplatform/shims/MultiplatformConcurrentHashSet.jvm.kt new file mode 100644 index 000000000000..423e8b97a81f --- /dev/null +++ b/fleet/multiplatform.shims/srcJvmMain/fleet/multiplatform/shims/MultiplatformConcurrentHashSet.jvm.kt @@ -0,0 +1,20 @@ +// 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 fleet.util.multiplatform.Actual +import java.util.concurrent.ConcurrentHashMap + +@Actual +internal fun MultiplatformConcurrentHashSetJvm(): MultiplatformConcurrentHashSet = MultiplatformConcurrentHashSetJvmImpl() + +private class MultiplatformConcurrentHashSetJvmImpl : MultiplatformConcurrentHashSet { + private val set = ConcurrentHashMap.newKeySet() + override val size: Int get() = set.size + override fun add(element: T): Boolean = set.add(element) + override fun iterator(): Iterator = set.iterator() + override fun addAll(elements: Collection): Boolean = set.addAll(elements) + override fun remove(element: T): Boolean = set.remove(element) + override fun contains(element: T): Boolean = set.contains(element) + override fun isEmpty(): Boolean = set.isEmpty() + override fun clear() = set.clear() +} diff --git a/fleet/multiplatform.shims/srcWasmJsMain/fleet/multiplatform/shims/ConcurrentHashSet.wasm.kt b/fleet/multiplatform.shims/srcWasmJsMain/fleet/multiplatform/shims/ConcurrentHashSet.wasm.kt deleted file mode 100644 index 4e4fdbe953ff..000000000000 --- a/fleet/multiplatform.shims/srcWasmJsMain/fleet/multiplatform/shims/ConcurrentHashSet.wasm.kt +++ /dev/null @@ -1,8 +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.multiplatform.shims - - -import fleet.util.multiplatform.Actual - -@Actual -internal fun ConcurrentHashSetWasmJs(): MutableSet = mutableSetOf() \ No newline at end of file diff --git a/fleet/multiplatform.shims/srcWasmJsMain/fleet/multiplatform/shims/MultiplatformConcurrentHashSet.wasm.kt b/fleet/multiplatform.shims/srcWasmJsMain/fleet/multiplatform/shims/MultiplatformConcurrentHashSet.wasm.kt new file mode 100644 index 000000000000..cd7dcce08574 --- /dev/null +++ b/fleet/multiplatform.shims/srcWasmJsMain/fleet/multiplatform/shims/MultiplatformConcurrentHashSet.wasm.kt @@ -0,0 +1,19 @@ +// 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 fleet.util.multiplatform.Actual + +@Actual +internal fun MultiplatformConcurrentHashSetWasmJs(): MultiplatformConcurrentHashSet = MultiplatformConcurrentHashSetWasmJs() + +private class MultiplatformConcurrentHashSetWasmJsImpl : MultiplatformConcurrentHashSet { + private val set = mutableSetOf() + override val size: Int get() = set.size + override fun add(element: T): Boolean = set.add(element) + override fun iterator(): Iterator = set.iterator() + override fun addAll(elements: Collection): Boolean = set.addAll(elements) + override fun remove(element: T): Boolean = set.remove(element) + override fun contains(element: T): Boolean = set.contains(element) + override fun isEmpty(): Boolean = set.isEmpty() + override fun clear() = set.clear() +} diff --git a/fleet/rpc.server/srcCommonMain/fleet/rpc/server/RpcExecutor.kt b/fleet/rpc.server/srcCommonMain/fleet/rpc/server/RpcExecutor.kt index 1412526a7bb6..6d2d5c7bc5ab 100644 --- a/fleet/rpc.server/srcCommonMain/fleet/rpc/server/RpcExecutor.kt +++ b/fleet/rpc.server/srcCommonMain/fleet/rpc/server/RpcExecutor.kt @@ -19,7 +19,7 @@ import kotlinx.coroutines.channels.SendChannel import kotlinx.coroutines.channels.consumeEach import kotlinx.serialization.json.Json import fleet.multiplatform.shims.ConcurrentHashMap -import fleet.multiplatform.shims.ConcurrentHashSet +import fleet.multiplatform.shims.MultiplatformConcurrentHashSet import fleet.util.async.withSupervisor import kotlinx.serialization.builtins.serializer import kotlin.coroutines.EmptyCoroutineContext @@ -101,9 +101,9 @@ class RpcExecutor private constructor( } private val requestJobs = ConcurrentHashMap() - private val routeRequests = ConcurrentHashMap>() + private val routeRequests = ConcurrentHashMap>() private val channels = ConcurrentHashMap() - private val routeChannels = ConcurrentHashMap>() + private val routeChannels = ConcurrentHashMap>() private suspend fun send(message: TransportMessage) { sendSuspend(::sendAsync, message) @@ -336,7 +336,7 @@ class RpcExecutor private constructor( require(previous == null) { "There is no way you can use the same channel twice ${descriptor.displayName}" } - routeChannels.computeIfAbsent(route) { ConcurrentHashSet() }.add(descriptor.uid) + routeChannels.computeIfAbsent(route) { MultiplatformConcurrentHashSet() }.add(descriptor.uid) return registeredStream } @@ -346,7 +346,7 @@ class RpcExecutor private constructor( route: UID?, ) { requestJobs[requestId] = requestJob - if (route != null) routeRequests.computeIfAbsent(route) { ConcurrentHashSet() }.add(requestId) + if (route != null) routeRequests.computeIfAbsent(route) { MultiplatformConcurrentHashSet() }.add(requestId) } private fun removeRequest(requestId: UID, route: UID?, jobAction: CompletableJob.() -> Unit) { diff --git a/fleet/util/core/srcCommonMain/fleet/util/DoOnce.kt b/fleet/util/core/srcCommonMain/fleet/util/DoOnce.kt index d0e79f1f92b1..832a98bfae26 100644 --- a/fleet/util/core/srcCommonMain/fleet/util/DoOnce.kt +++ b/fleet/util/core/srcCommonMain/fleet/util/DoOnce.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.util -import fleet.multiplatform.shims.ConcurrentHashSet +import fleet.multiplatform.shims.MultiplatformConcurrentHashSet import kotlin.jvm.JvmInline sealed class DoOnce { @@ -10,7 +10,7 @@ sealed class DoOnce { @JvmInline value class Id(val id: String) - private val done = ConcurrentHashSet() + private val done = MultiplatformConcurrentHashSet() fun doOnce(id: String, body: () -> T): T? { return if (done.add(id)) body() else null diff --git a/fleet/util/core/srcJvmMain/fleet/util/NetUtils.jvm.kt b/fleet/util/core/srcJvmMain/fleet/util/NetUtils.jvm.kt index c49f52e37cc0..84a396ded558 100644 --- a/fleet/util/core/srcJvmMain/fleet/util/NetUtils.jvm.kt +++ b/fleet/util/core/srcJvmMain/fleet/util/NetUtils.jvm.kt @@ -1,13 +1,13 @@ // 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.multiplatform.shims.ConcurrentHashSet +import fleet.multiplatform.shims.MultiplatformConcurrentHashSet import java.net.InetAddress import java.net.ServerSocket import kotlin.random.Random import kotlin.random.nextInt -private val used = ConcurrentHashSet() +private val used = MultiplatformConcurrentHashSet() // https://en.wikipedia.org/wiki/Ephemeral_port private val portsRange = (1 shl 15) + (1 shl 14) until (1 shl 16)