Fix NegativeArraySizeException crash in relay subscriptions

ProfileSubscription and NotificationSubscription kept their subscription
state in plain LinkedHashMaps that are mutated from UI/coroutine threads
(updateFilter, closeSub, updateFilters) while relay I/O threads iterate
them in onIncomingMessage. Iterating a LinkedHashMap's values/entries
while another thread structurally modifies it crashes with
NegativeArraySizeException on Android, e.g. from
ProfileSubscription.updateFilters via checkForNewRelaysAndUpdateAllFilters.

- Use ConcurrentHashMap for all shared maps in both subscriptions and
  concurrent key sets for the per-subscription relay sets
- Replace containsKey + get(!!) with an atomic getOrPut in
  NotificationSubscription.updateFilter

Add unit tests: functional coverage for the subscribe/unsubscribe
lifecycle and concurrency regression tests that reliably reproduce the
race on the pre-fix implementation.
This commit is contained in:
greenart7c3
2026-07-29 07:06:07 -03:00
parent 932434440d
commit 1056fb39db
5 changed files with 457 additions and 16 deletions
@@ -36,6 +36,7 @@ import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter
import com.vitorpamplona.quartz.nip46RemoteSigner.NostrConnectEvent
import com.vitorpamplona.quartz.utils.TimeUtils
import java.util.UUID
import java.util.concurrent.ConcurrentHashMap
import kotlinx.coroutines.launch
class NotificationSubscription(
@@ -43,7 +44,10 @@ class NotificationSubscription(
val appContext: Context,
) : RelayConnectionListener {
private val eventNotificationConsumer = EventNotificationConsumer(appContext)
private val subIds = mutableMapOf<String, String>()
// Concurrent: mutated in updateFilter/closeAllSubs (coroutines) while relay I/O
// threads iterate it in onIncomingMessage (containsValue).
private val subIds = ConcurrentHashMap<String, String>()
init {
// listens until the app crashes.
@@ -101,11 +105,9 @@ class NotificationSubscription(
if (connRelays.isEmpty()) continue
activeSubKeys.add(subKey)
if (!subIds.containsKey(subKey)) {
subIds[subKey] = UUID.randomUUID().toString()
}
val subId = subIds.getOrPut(subKey) { UUID.randomUUID().toString() }
client.subscribe(
subIds[subKey]!!,
subId,
connRelays.associateWith {
listOf(
Filter(
@@ -125,11 +127,9 @@ class NotificationSubscription(
.filterNot { RelayHealthTracker.isDead(it) }
if (relays.isNotEmpty()) {
activeSubKeys.add(account.hexKey)
if (!subIds.containsKey(account.hexKey)) {
subIds[account.hexKey] = UUID.randomUUID().toString()
}
val subId = subIds.getOrPut(account.hexKey) { UUID.randomUUID().toString() }
client.subscribe(
subIds[account.hexKey]!!,
subId,
relays.associateWith {
listOf(
Filter(
@@ -39,6 +39,7 @@ import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
import com.vitorpamplona.quartz.nip65RelayList.AdvertisedRelayListEvent
import com.vitorpamplona.quartz.utils.TimeUtils
import java.util.UUID
import java.util.concurrent.ConcurrentHashMap
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Job
import kotlinx.coroutines.delay
@@ -52,17 +53,21 @@ class ProfileSubscription(
val appContext: Context,
val scope: CoroutineScope,
) : RelayConnectionListener {
// All maps are concurrent: mutated from UI/coroutine threads (updateFilter, closeSub)
// and read/iterated from relay I/O threads (onIncomingMessage) at the same time.
// Plain LinkedHashMaps crash with NegativeArraySizeException under that race.
// hexKey -> subId of the kind-0 metadata subscription
private val subIds = mutableMapOf<String, String>()
private val subIds = ConcurrentHashMap<String, String>()
// hexKey -> subId of the kind-10002 relay list subscription that runs before the metadata one
private val relayListSubIds = mutableMapOf<String, String>()
private val relaysPerSubId = mutableMapOf<String, MutableSet<NormalizedRelayUrl>>()
private val timeoutJobs = mutableMapOf<String, Job>()
private val relayListSubIds = ConcurrentHashMap<String, String>()
private val relaysPerSubId = ConcurrentHashMap<String, MutableSet<NormalizedRelayUrl>>()
private val timeoutJobs = ConcurrentHashMap<String, Job>()
// hexKey -> the cached Account instance a live composable cares about. Updating the
// StateFlows on these instances is what reflects fresh metadata in the UI.
private val accounts = mutableMapOf<String, Account>()
private val accounts = ConcurrentHashMap<String, Account>()
init {
// listens until the app crashes.
@@ -162,7 +167,7 @@ class ProfileSubscription(
val subId = relayListSubIds.getOrPut(account.hexKey) { UUID.randomUUID().toString() }
val relayListFilter = createRelayListFilter(account)
timeoutJobs.remove(subId)?.cancel()
relaysPerSubId[subId] = relayListFilter.keys.toMutableSet()
relaysPerSubId[subId] = ConcurrentHashMap.newKeySet<NormalizedRelayUrl>().apply { addAll(relayListFilter.keys) }
client.subscribe(subId, relayListFilter)
timeoutJobs[subId] = scope.launch {
delay(EOSE_TIMEOUT_MS)
@@ -178,7 +183,7 @@ class ProfileSubscription(
val subId = subIds.getOrPut(account.hexKey) { UUID.randomUUID().toString() }
val profileFilter = createProfileFilter(account)
timeoutJobs.remove(subId)?.cancel()
relaysPerSubId[subId] = profileFilter.keys.toMutableSet()
relaysPerSubId[subId] = ConcurrentHashMap.newKeySet<NormalizedRelayUrl>().apply { addAll(profileFilter.keys) }
client.subscribe(subId, profileFilter)
timeoutJobs[subId] = scope.launch {
delay(EOSE_TIMEOUT_MS)
@@ -0,0 +1,171 @@
package com.greenart7c3.nostrsigner.service
import androidx.collection.LruCache
import com.greenart7c3.nostrsigner.Amber
import com.greenart7c3.nostrsigner.BuildFlavorChecker
import com.greenart7c3.nostrsigner.LocalPreferences
import com.greenart7c3.nostrsigner.database.AppDatabase
import com.greenart7c3.nostrsigner.database.ApplicationDao
import com.greenart7c3.nostrsigner.database.ApplicationEntity
import com.greenart7c3.nostrsigner.models.Account
import com.vitorpamplona.quartz.nip01Core.core.Event
import com.vitorpamplona.quartz.nip01Core.relay.client.NostrClient
import com.vitorpamplona.quartz.nip01Core.relay.client.single.IRelayClient
import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.EventMessage
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
import io.mockk.coEvery
import io.mockk.every
import io.mockk.mockk
import io.mockk.mockkObject
import io.mockk.unmockkObject
import io.mockk.verify
import java.util.Collections
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.SupervisorJob
import kotlinx.coroutines.cancel
import kotlinx.coroutines.runBlocking
import org.junit.After
import org.junit.Assert.assertEquals
import org.junit.Assert.assertTrue
import org.junit.Assume.assumeFalse
import org.junit.Before
import org.junit.Test
class NotificationSubscriptionTest {
private val scope = CoroutineScope(SupervisorJob() + Dispatchers.Default)
private val sentSubIds = Collections.synchronizedList(mutableListOf<String>())
private lateinit var client: NostrClient
private lateinit var dao: ApplicationDao
private lateinit var account: Account
private lateinit var subscription: NotificationSubscription
@Before
fun setUp() {
client = mockk(relaxed = true)
every { client.subscribe(capture(sentSubIds), any()) } returns Unit
dao = mockk()
val database = mockk<AppDatabase>()
every { database.dao() } returns dao
val amber = mockk<Amber>(relaxed = true)
every { amber.getDatabase(any()) } returns database
every { amber.notificationCache } returns LruCache(512)
installAmberInstance(amber)
mockkObject(LocalPreferences)
account = newTestAccount(scope)
coEvery { LocalPreferences.allAccounts(any()) } returns listOf(account)
subscription = NotificationSubscription(client, mockk(relaxed = true))
}
@After
fun tearDown() {
unmockkObject(LocalPreferences)
scope.cancel()
}
private fun connection() = ApplicationEntity.empty().copy(
key = "conn1",
relays = listOf(NormalizedRelayUrl("wss://relay1.example.com")),
localKey = "ab".repeat(32),
)
@Test
fun `updateFilter subscribes once per connection with a local key`() = runBlocking {
// updateFilter is a no-op on the offline flavor by design
assumeFalse(BuildFlavorChecker.isOfflineFlavor())
coEvery { dao.getAll(account.hexKey) } returns listOf(connection())
every { dao.getAllRelayLists() } returns emptyList()
subscription.updateFilter()
assertEquals(1, sentSubIds.size)
}
@Test
fun `updateFilter reuses the same subscription id on refresh`() = runBlocking {
assumeFalse(BuildFlavorChecker.isOfflineFlavor())
coEvery { dao.getAll(account.hexKey) } returns listOf(connection())
every { dao.getAllRelayLists() } returns emptyList()
subscription.updateFilter()
subscription.updateFilter()
assertEquals(2, sentSubIds.size)
assertEquals(sentSubIds[0], sentSubIds[1])
}
@Test
fun `updateFilter unsubscribes connections that disappeared`() = runBlocking {
assumeFalse(BuildFlavorChecker.isOfflineFlavor())
coEvery { dao.getAll(account.hexKey) } returns listOf(connection())
every { dao.getAllRelayLists() } returns emptyList()
subscription.updateFilter()
val subId = sentSubIds.single()
coEvery { dao.getAll(account.hexKey) } returns emptyList()
subscription.updateFilter()
verify(exactly = 1) { client.unsubscribe(subId) }
}
/**
* Regression test for the NegativeArraySizeException class of bugs: onIncomingMessage
* iterates subIds (containsValue) on relay I/O threads while updateFilter/closeAllSubs
* mutate it from coroutines.
*/
@Test
fun `updateFilter closeAllSubs and onIncomingMessage do not throw when run concurrently`() {
coEvery { dao.getAll(account.hexKey) } returns listOf(connection())
every { dao.getAllRelayLists() } returns emptyList()
val relay = mockk<IRelayClient>()
val event = mockk<Event>(relaxed = true)
val errors = Collections.synchronizedList(mutableListOf<Throwable>())
val threads = mutableListOf<Thread>()
repeat(2) {
threads += Thread {
repeat(100) {
try {
runBlocking { subscription.updateFilter() }
} catch (e: Throwable) {
errors.add(e)
}
}
}
}
threads += Thread {
repeat(100) {
try {
subscription.closeAllSubs()
} catch (e: Throwable) {
errors.add(e)
}
}
}
repeat(2) {
threads += Thread {
repeat(300) { i ->
try {
subscription.onIncomingMessage(relay, "", EventMessage("unknown-sub-$i", event))
} catch (e: Throwable) {
errors.add(e)
}
}
}
}
threads.forEach { it.start() }
threads.forEach { it.join() }
assertTrue(
"Expected no exceptions, got: ${errors.firstOrNull()?.stackTraceToString()}",
errors.isEmpty(),
)
}
}
@@ -0,0 +1,223 @@
package com.greenart7c3.nostrsigner.service
import com.greenart7c3.nostrsigner.Amber
import com.greenart7c3.nostrsigner.BuildFlavorChecker
import com.greenart7c3.nostrsigner.LocalPreferences
import com.greenart7c3.nostrsigner.models.Account
import com.greenart7c3.nostrsigner.models.AmberSettings
import com.greenart7c3.nostrsigner.models.ProfileFetchInterval
import com.vitorpamplona.quartz.nip01Core.core.Event
import com.vitorpamplona.quartz.nip01Core.relay.client.NostrClient
import com.vitorpamplona.quartz.nip01Core.relay.client.single.IRelayClient
import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.EoseMessage
import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.EventMessage
import com.vitorpamplona.quartz.utils.TimeUtils
import io.mockk.Runs
import io.mockk.every
import io.mockk.just
import io.mockk.mockk
import io.mockk.mockkObject
import io.mockk.unmockkObject
import io.mockk.verify
import java.util.Collections
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.SupervisorJob
import kotlinx.coroutines.cancel
import kotlinx.coroutines.runBlocking
import org.junit.After
import org.junit.Assert.assertTrue
import org.junit.Assume.assumeFalse
import org.junit.Before
import org.junit.Test
class ProfileSubscriptionTest {
private val scope = CoroutineScope(SupervisorJob() + Dispatchers.Default)
private val sentSubIds = Collections.synchronizedList(mutableListOf<String>())
private lateinit var client: NostrClient
private lateinit var subscription: ProfileSubscription
@Before
fun setUp() {
client = mockk(relaxed = true)
every { client.subscribe(capture(sentSubIds), any()) } returns Unit
installAmber(ProfileFetchInterval.ALWAYS)
mockkObject(LocalPreferences)
every { LocalPreferences.loadSettingsFromEncryptedStorage(any()) } returns AmberSettings()
every { LocalPreferences.getUserRelays(any(), any()) } returns emptyList()
every { LocalPreferences.getLastMetadataUpdate(any(), any()) } returns 0L
every { LocalPreferences.getLastCheck(any(), any()) } returns 0L
every { LocalPreferences.setLastCheck(any(), any(), any()) } just Runs
every { LocalPreferences.setLastMetadataUpdate(any(), any(), any()) } just Runs
subscription = ProfileSubscription(client, mockk(relaxed = true), scope)
}
@After
fun tearDown() {
unmockkObject(LocalPreferences)
scope.cancel()
}
private fun installAmber(interval: ProfileFetchInterval) {
val amber = mockk<Amber>(relaxed = true)
every { amber.settings } returns AmberSettings(profileFetchInterval = interval)
installAmberInstance(amber)
}
@Test
fun `updateFilter subscribes to the user relay list when interval is ALWAYS`() = runBlocking {
// updateFilter is a no-op on the offline flavor by design
assumeFalse(BuildFlavorChecker.isOfflineFlavor())
subscription.updateFilter(newTestAccount(scope))
verify(exactly = 1) { client.subscribe(any(), any()) }
}
@Test
fun `updateFilter does not subscribe when fetch interval is NEVER`() = runBlocking {
assumeFalse(BuildFlavorChecker.isOfflineFlavor())
installAmber(ProfileFetchInterval.NEVER)
subscription.updateFilter(newTestAccount(scope))
verify(exactly = 0) { client.subscribe(any(), any()) }
}
@Test
fun `closeSub unsubscribes the account subscriptions`() = runBlocking {
assumeFalse(BuildFlavorChecker.isOfflineFlavor())
val account = newTestAccount(scope)
subscription.updateFilter(account)
subscription.closeSub(account)
verify(exactly = 1) { client.unsubscribe(any()) }
}
/**
* Regression test for the NegativeArraySizeException crash: updateFilters iterated
* accounts.values while the UI thread added/removed accounts on a plain LinkedHashMap.
* Uses a large tracked set and a throttled interval (no relay work) so the
* iteration-vs-mutation race window is hit thousands of times.
*/
@Test
fun `updateFilters does not throw when accounts are added and removed concurrently`() {
installAmber(ProfileFetchInterval.ONE_HOUR)
every { LocalPreferences.getLastCheck(any(), any()) } returns TimeUtils.now()
val tracked = Collections.synchronizedList(mutableListOf<Account>())
runBlocking {
repeat(200) {
val account = newTestAccount(scope)
tracked.add(account)
subscription.updateFilter(account)
}
}
val errors = Collections.synchronizedList(mutableListOf<Throwable>())
val threads = mutableListOf<Thread>()
repeat(3) {
threads += Thread {
repeat(100) {
try {
runBlocking { subscription.updateFilters() }
} catch (e: Throwable) {
errors.add(e)
}
}
}
}
// Churn threads recycle the existing pool: closeSub + updateFilter on the same
// accounts gives endless structural map mutations without allocating new objects.
repeat(2) {
threads += Thread {
repeat(400) {
try {
tracked.randomOrNull()?.let { account ->
subscription.closeSub(account)
runBlocking { subscription.updateFilter(account) }
}
} catch (e: Throwable) {
errors.add(e)
}
}
}
}
repeat(2) {
threads += Thread {
repeat(400) {
try {
tracked.randomOrNull()?.let { subscription.closeSub(it) }
} catch (e: Throwable) {
errors.add(e)
}
}
}
}
threads.forEach { it.start() }
threads.forEach { it.join() }
assertTrue(
"Expected no exceptions, got: ${errors.firstOrNull()?.stackTraceToString()}",
errors.isEmpty(),
)
}
/**
* EOSE handling looks up subscription ids in the same maps that updateFilter/closeSub
* mutate; these run on relay I/O threads and coroutines at the same time.
*/
@Test
fun `onIncomingMessage does not throw when subscriptions change concurrently`() {
val relay = mockk<IRelayClient>()
every { relay.url } returns AmberSettings().defaultProfileRelays.first()
val event = mockk<Event>()
every { event.kind } returns -1
val accounts = (1..8).map { newTestAccount(scope) }
runBlocking { accounts.forEach { subscription.updateFilter(it) } }
val errors = Collections.synchronizedList(mutableListOf<Throwable>())
val threads = mutableListOf<Thread>()
repeat(3) {
threads += Thread {
repeat(200) {
try {
val subId = sentSubIds.toList().randomOrNull() ?: return@repeat
subscription.onIncomingMessage(relay, "", EoseMessage(subId))
subscription.onIncomingMessage(relay, "", EventMessage(subId, event))
} catch (e: Throwable) {
errors.add(e)
}
}
}
}
threads += Thread {
repeat(100) {
try {
runBlocking { subscription.updateFilters() }
} catch (e: Throwable) {
errors.add(e)
}
}
}
threads += Thread {
repeat(50) {
try {
subscription.closeSub(accounts.random())
} catch (e: Throwable) {
errors.add(e)
}
}
}
threads.forEach { it.start() }
threads.forEach { it.join() }
assertTrue(
"Expected no exceptions, got: ${errors.firstOrNull()?.stackTraceToString()}",
errors.isEmpty(),
)
}
}
@@ -0,0 +1,42 @@
package com.greenart7c3.nostrsigner.service
import com.greenart7c3.nostrsigner.Amber
import com.greenart7c3.nostrsigner.models.Account
import com.vitorpamplona.quartz.nip01Core.signers.NostrSignerInternal
import io.mockk.mockk
import java.util.concurrent.atomic.AtomicInteger
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.flow.MutableStateFlow
private val accountIds = AtomicInteger(0)
// One shared relaxed mock for all test accounts: each mockk() generates proxy
// classes, and hundreds of them blow the test heap. The signer is never exercised.
private val sharedSigner = mockk<NostrSignerInternal>(relaxed = true)
/**
* Installs a mock [Amber] into the `Amber.instance` lateinit field, which has a
* private setter and is normally only assigned by the Application onCreate.
* The backing field of a companion `lateinit var` is a private static field on
* the outer class.
*/
fun installAmberInstance(amber: Amber) {
val field = Amber::class.java.getDeclaredField("instance")
field.isAccessible = true
field.set(null, amber)
}
/** Builds a minimal [Account] with a unique key; the signer itself is never exercised. */
fun newTestAccount(scope: CoroutineScope): Account {
val id = accountIds.incrementAndGet()
return Account(
signer = sharedSigner,
hexKey = "%064x".format(id),
npub = "npub$id",
name = MutableStateFlow(""),
picture = MutableStateFlow(""),
signPolicy = 0,
didBackup = false,
scope = scope,
)
}