[fleet] rename multiplatform ConcurrentHashSet and drop inheritance from the MutableSet

Renaming will prevent misuse of multiplatform collection in non-fleet projects.
 Dropping an inheritance allows us to avoid tricky bugs related to the extension functions on a MutableSet that become accidentally available on the concurrentSet, while not being really concurrent

GitOrigin-RevId: a8dca907c60d772d89674553814320f4c5fd2032
This commit is contained in:
Alexander Zolotov
2025-11-08 18:46:14 +00:00
committed by intellij-monorepo-bot
parent 9c60bae14e
commit dedc339a98
12 changed files with 110 additions and 39 deletions
@@ -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 <T : Any> SubscriptionScope.distinct(producer: Producer<T>): Producer<T> =
run {
val broadcast = Broadcaster<T>()
val memory = HashMap<T, MutableSet<Match<T>>>()
val memory = HashMap<T, MultiplatformConcurrentHashSet<Match<T>>>()
producer.collect { token ->
val v = token.match.value
when (token.added) {
true -> {
when (val matches = memory[v]) {
null -> {
val ms = ConcurrentHashSet<Match<T>>()
val ms = MultiplatformConcurrentHashSet<Match<T>>()
memory[v] = ms
ms.add(token.match)
broadcast(Token(true, distinctMatchValue(v, ms)))
@@ -29,7 +29,7 @@ internal fun <T : Any> SubscriptionScope.distinct(producer: Producer<T>): 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 <T : Any> SubscriptionScope.distinct(producer: Producer<T>): Produc
}
}
private fun <T> distinctMatchValue(v: T, ms: Set<Match<T>>): Match<T> =
private fun <T> distinctMatchValue(v: T, ms: MultiplatformConcurrentHashSet<Match<T>>): Match<T> =
Match.validatable(v) {
val anyValid = ms.any { match ->
match.validate() == ValidationResultEnum.Valid
@@ -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 <K> ConcurrentHashSet(): MutableSet<K> = linkToActual()
@@ -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<Nothing> {
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()
}
@@ -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 <T> MultiplatformConcurrentHashSet(): MultiplatformConcurrentHashSet<T> = linkToActual()
interface MultiplatformConcurrentHashSet<out T> : Iterable<T> {
companion object {
fun <T> empty(): MultiplatformConcurrentHashSet<T> {
return EmptySet as MultiplatformConcurrentHashSet<T>
}
}
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<Nothing> {
override val size: Int get() = 0
override fun add(element: Nothing): Boolean = throw UnsupportedOperationException()
override fun addAll(elements: Collection<Nothing>): 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<Nothing> = EmptyIterator
override fun equals(other: Any?): Boolean = other is MultiplatformConcurrentHashSet<*> && other.isEmpty()
override fun hashCode(): Int = 0
}
@@ -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 <K> ConcurrentHashSetJvm(): MutableSet<K> = JavaConcurrentHashMap.newKeySet()
@@ -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<K, V>.newKeySet()",
"java.util.concurrent.ConcurrentHashMap"))
internal fun <K> ConcurrentHashSet(): MutableSet<K> = JavaConcurrentHashMap.newKeySet()
@@ -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 <K> MultiplatformConcurrentHashSetJvm(): MultiplatformConcurrentHashSet<K> = MultiplatformConcurrentHashSetJvmImpl()
private class MultiplatformConcurrentHashSetJvmImpl<T> : MultiplatformConcurrentHashSet<T> {
private val set = ConcurrentHashMap.newKeySet<T>()
override val size: Int get() = set.size
override fun add(element: T): Boolean = set.add(element)
override fun iterator(): Iterator<T> = set.iterator()
override fun addAll(elements: Collection<T>): 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()
}
@@ -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 <K> ConcurrentHashSetWasmJs(): MutableSet<K> = mutableSetOf<K>()
@@ -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 <K> MultiplatformConcurrentHashSetWasmJs(): MultiplatformConcurrentHashSet<K> = MultiplatformConcurrentHashSetWasmJs()
private class MultiplatformConcurrentHashSetWasmJsImpl<T> : MultiplatformConcurrentHashSet<T> {
private val set = mutableSetOf<T>()
override val size: Int get() = set.size
override fun add(element: T): Boolean = set.add(element)
override fun iterator(): Iterator<T> = set.iterator()
override fun addAll(elements: Collection<T>): 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()
}
@@ -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<UID, CompletableJob>()
private val routeRequests = ConcurrentHashMap<UID/*route*/, MutableSet<UID/*requestId*/>>()
private val routeRequests = ConcurrentHashMap<UID/*route*/, MultiplatformConcurrentHashSet<UID/*requestId*/>>()
private val channels = ConcurrentHashMap<UID, InternalStreamDescriptor>()
private val routeChannels = ConcurrentHashMap<UID/*route*/, MutableSet<UID/*channelId*/>>()
private val routeChannels = ConcurrentHashMap<UID/*route*/, MultiplatformConcurrentHashSet<UID/*channelId*/>>()
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) {
@@ -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<String>()
private val done = MultiplatformConcurrentHashSet<String>()
fun <T> doOnce(id: String, body: () -> T): T? {
return if (done.add(id)) body() else null
@@ -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<Int>()
private val used = MultiplatformConcurrentHashSet<Int>()
// https://en.wikipedia.org/wiki/Ephemeral_port
private val portsRange = (1 shl 15) + (1 shl 14) until (1 shl 16)