Merge pull request #3666 from davotoula/fix/addressable-author-loader-cme-alternative

Stop ConcurrentModificationException in AddressableAuthorRelayLoaderSubAssembler
This commit is contained in:
Vitor Pamplona
2026-07-22 12:43:16 -04:00
committed by GitHub
4 changed files with 225 additions and 10 deletions
+1 -1
View File
@@ -8,6 +8,6 @@
</component>
<component name="KotlinJpsPluginSettings">
<option name="externalSystemId" value="Gradle" />
<option name="version" value="2.4.0" />
<option name="version" value="2.4.10" />
</component>
</project>
@@ -21,11 +21,14 @@
package com.vitorpamplona.amethyst.service.relayClient.reqCommand.event.loaders
import com.vitorpamplona.amethyst.commons.relayClient.eoseManagers.IEoseManager
import com.vitorpamplona.amethyst.commons.service.BundledUpdate
import com.vitorpamplona.amethyst.model.AddressableNote
import com.vitorpamplona.amethyst.model.LocalCache
import com.vitorpamplona.amethyst.service.relayClient.reqCommand.event.EventFinderQueryState
import com.vitorpamplona.amethyst.service.relayClient.reqCommand.user.UserFinderFilterAssembler
import com.vitorpamplona.amethyst.service.relayClient.reqCommand.user.UserFinderQueryState
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.IO
/**
* Bridges missing-addressable-note authors into [UserFinderFilterAssembler].
@@ -42,9 +45,24 @@ class AddressableAuthorRelayLoaderSubAssembler(
val allKeys: () -> Set<EventFinderQueryState>,
val userFinder: UserFinderFilterAssembler,
) : IEoseManager {
private val activeSubscriptions = mutableSetOf<UserFinderQueryState>()
// Private monitor: @Synchronized locks on `this`, which leaves the instance's monitor
// reachable to anything holding a reference to this assembler.
private val lock = Any()
// Only ever touched while holding [lock]. See commit() and destroy().
private var activeSubscriptions: Set<UserFinderQueryState> = emptySet()
private var destroyed = false
// Keeps the scan off the caller's thread. invalidateFilters() is reached synchronously from
// ComposeSubscriptionManager.subscribe/unsubscribe on every note composable mount/unmount,
// and those are documented "called by main. Keep it really fast."
private val bundler = BundledUpdate(500, Dispatchers.IO)
override fun invalidateFilters(ignoreIfDoing: Boolean) {
bundler.invalidate(ignoreIfDoing, ::forceInvalidate)
}
private fun forceInvalidate() {
val needed = mutableSetOf<UserFinderQueryState>()
allKeys().forEach { key ->
@@ -57,18 +75,39 @@ class AddressableAuthorRelayLoaderSubAssembler(
}
}
val toAdd = needed - activeSubscriptions
val toRemove = activeSubscriptions - needed
commit(needed)
}
userFinder.subscribe(toAdd.toList())
userFinder.unsubscribe(toRemove.toList())
/**
* Serializes against [destroy] — the one caller the bundler cannot order, because
* `bundler.cancel()` cannot stop a body that is already running (it has no suspension points).
*
* The scan in [forceInvalidate] stays outside [lock], so [destroy] never waits on a
* [LocalCache] sweep. It can still wait on the two calls below, which are bounded: a pair of
* map updates inside [UserFinderFilterAssembler] plus the coroutine launches its
* `invalidateKeys()` fans out to.
*
* Calling [userFinder] while holding [lock] relies on subscribe/unsubscribe only taking
* ComposeSubscriptionManager's own lock and deferring real work to bundled coroutines — they
* never call back into this class. Revisit if that changes.
*/
private fun commit(needed: Set<UserFinderQueryState>) {
synchronized(lock) {
if (destroyed) return
activeSubscriptions.clear()
activeSubscriptions.addAll(needed)
userFinder.subscribe((needed - activeSubscriptions).toList())
userFinder.unsubscribe((activeSubscriptions - needed).toList())
activeSubscriptions = needed
}
}
override fun destroy() {
userFinder.unsubscribe(activeSubscriptions.toList())
activeSubscriptions.clear()
synchronized(lock) {
destroyed = true
bundler.cancel()
userFinder.unsubscribe(activeSubscriptions.toList())
activeSubscriptions = emptySet()
}
}
}
@@ -0,0 +1,170 @@
/*
* Copyright (c) 2025 Vitor Pamplona
*
* Permission is hereby granted, free of charge, to any person obtaining a copy of
* this software and associated documentation files (the "Software"), to deal in
* the Software without restriction, including without limitation the rights to use,
* copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the
* Software, and to permit persons to whom the Software is furnished to do so,
* subject to the following conditions:
*
* The above copyright notice and this permission notice shall be included in all
* copies or substantial portions of the Software.
*
* THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
* IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS
* FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR
* COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN
* AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION
* WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE.
*/
package com.vitorpamplona.amethyst.service.relayClient.reqCommand.event.loaders
import com.vitorpamplona.amethyst.model.Account
import com.vitorpamplona.amethyst.model.LocalCache
import com.vitorpamplona.amethyst.service.relayClient.reqCommand.event.EventFinderQueryState
import com.vitorpamplona.amethyst.service.relayClient.reqCommand.user.UserFinderFilterAssembler
import com.vitorpamplona.amethyst.service.relayClient.reqCommand.user.UserFinderQueryState
import com.vitorpamplona.quartz.nip01Core.core.Address
import io.mockk.every
import io.mockk.mockk
import io.mockk.verify
import org.junit.Assert.assertEquals
import org.junit.Assert.assertFalse
import org.junit.Assert.assertTrue
import org.junit.Test
import java.util.concurrent.CountDownLatch
import java.util.concurrent.TimeUnit
import java.util.concurrent.atomic.AtomicInteger
import kotlin.concurrent.thread
class AddressableAuthorRelayLoaderSubAssemblerTest {
/**
* Unique per-test-class identities so the shared [LocalCache] singleton isn't polluted with
* notes another test class also claims.
*/
private fun stubKeys(count: Int): Set<EventFinderQueryState> {
val account = mockk<Account>()
return (1..count).mapTo(mutableSetOf()) { i ->
val address = Address(30023, "ad04%060x".format(i), "d$i")
EventFinderQueryState(LocalCache.getOrCreateAddressableNoteInternal(address), account)
}
}
/**
* Regression test for the `ConcurrentModificationException` in `SetsKt.minus` reported from
* `DefaultDispatcher-worker-70`: this manager is genuinely re-entered from several threads at
* once, because `ComposeSubscriptionManager.subscribe`/`unsubscribe` call `invalidateKeys()`
* *after* releasing their own lock, and `LifecycleAwareSubscription`'s 30s grace-period
* unsubscribe fires on a `Dispatchers.Default` worker.
*
* Overlap is detected through the injected `allKeys()` lambda rather than by catching — once
* the body is bundled its throwables are swallowed by `BundledUpdate`'s
* `CoroutineExceptionHandler` and would never reach the test thread.
*/
@Test
fun concurrentInvalidateFiltersNeverOverlap() {
val userFinder = mockk<UserFinderFilterAssembler>(relaxed = true)
val keys = stubKeys(50)
val inFlight = AtomicInteger(0)
val overlaps = AtomicInteger(0)
val assembler =
AddressableAuthorRelayLoaderSubAssembler(
LocalCache,
{
if (inFlight.incrementAndGet() > 1) overlaps.incrementAndGet()
try {
keys
} finally {
inFlight.decrementAndGet()
}
},
userFinder,
)
val start = CountDownLatch(1)
try {
val threads =
(1..8).map {
thread(start = false) {
start.await()
repeat(500) { assembler.invalidateFilters() }
}
}
threads.forEach { it.start() }
start.countDown()
threads.forEach { it.join() }
} finally {
assembler.destroy()
}
assertEquals("forceInvalidate bodies overlapped", 0, overlaps.get())
}
/** The manager still does its job: unresolved stub authors reach the user finder. */
@Test
fun bridgesMissingAuthorsIntoUserFinder() {
val userFinder = mockk<UserFinderFilterAssembler>(relaxed = true)
val keys = stubKeys(3)
val assembler = AddressableAuthorRelayLoaderSubAssembler(LocalCache, { keys }, userFinder)
try {
assembler.invalidateFilters()
verify(timeout = 3000) {
userFinder.subscribe(
match<List<UserFinderQueryState>> { it.size == 3 },
)
}
} finally {
assembler.destroy()
}
}
/**
* `destroy()` must win against a body that is already past its `allKeys()` scan: otherwise the
* body re-subscribes authors `destroy()` just released, leaving live kind-0/10002 REQs (and
* retained `User`/`Account` references) for a dead account after logout.
*
* The body is gated inside the injected `allKeys()` lambda so `destroy()` provably runs
* underneath an in-flight run rather than racing it by luck.
*/
@Test
fun destroyDuringInFlightInvalidateDoesNotLeakSubscriptions() {
val userFinder = mockk<UserFinderFilterAssembler>(relaxed = true)
val subscribeHappened = CountDownLatch(1)
every { userFinder.subscribe(any<List<UserFinderQueryState>>()) } answers {
subscribeHappened.countDown()
}
val keys = stubKeys(3)
val invalidateEntered = CountDownLatch(1)
val destroyFinished = CountDownLatch(1)
val assembler =
AddressableAuthorRelayLoaderSubAssembler(
LocalCache,
{
invalidateEntered.countDown()
destroyFinished.await(5, TimeUnit.SECONDS)
keys
},
userFinder,
)
try {
assembler.invalidateFilters()
assertTrue("bundled run never started", invalidateEntered.await(5, TimeUnit.SECONDS))
} finally {
assembler.destroy()
destroyFinished.countDown()
}
// The gated body resumes the instant destroyFinished counts down, so a leak shows up
// immediately; the wait only has to outlast that hand-off.
assertFalse(
"in-flight body subscribed after destroy()",
subscribeHappened.await(500, TimeUnit.MILLISECONDS),
)
}
}
@@ -21,6 +21,12 @@
package com.vitorpamplona.amethyst.commons.relayClient.eoseManagers
interface IEoseManager {
/**
* May be called from any thread, concurrently with itself and with [destroy], and is reached
* synchronously from main on every composable mount/unmount. Implementations must return fast
* and must serialize their own state — see [BaseEoseManager], which does both by routing the
* work through a [com.vitorpamplona.amethyst.commons.service.BundledUpdate].
*/
fun invalidateFilters(ignoreIfDoing: Boolean = false)
fun destroy()