From 1056fb39dbbc93b7750e04f91f2368e2e0a27769 Mon Sep 17 00:00:00 2001 From: greenart7c3 Date: Wed, 29 Jul 2026 07:06:07 -0300 Subject: [PATCH] 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. --- .../service/NotificationSubscription.kt | 18 +- .../service/ProfileSubscription.kt | 19 +- .../service/NotificationSubscriptionTest.kt | 171 ++++++++++++++ .../service/ProfileSubscriptionTest.kt | 223 ++++++++++++++++++ .../service/SubscriptionTestUtils.kt | 42 ++++ 5 files changed, 457 insertions(+), 16 deletions(-) create mode 100644 app/src/test/java/com/greenart7c3/nostrsigner/service/NotificationSubscriptionTest.kt create mode 100644 app/src/test/java/com/greenart7c3/nostrsigner/service/ProfileSubscriptionTest.kt create mode 100644 app/src/test/java/com/greenart7c3/nostrsigner/service/SubscriptionTestUtils.kt diff --git a/app/src/main/java/com/greenart7c3/nostrsigner/service/NotificationSubscription.kt b/app/src/main/java/com/greenart7c3/nostrsigner/service/NotificationSubscription.kt index 47c473f0..8a56146b 100644 --- a/app/src/main/java/com/greenart7c3/nostrsigner/service/NotificationSubscription.kt +++ b/app/src/main/java/com/greenart7c3/nostrsigner/service/NotificationSubscription.kt @@ -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() + + // Concurrent: mutated in updateFilter/closeAllSubs (coroutines) while relay I/O + // threads iterate it in onIncomingMessage (containsValue). + private val subIds = ConcurrentHashMap() 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( diff --git a/app/src/main/java/com/greenart7c3/nostrsigner/service/ProfileSubscription.kt b/app/src/main/java/com/greenart7c3/nostrsigner/service/ProfileSubscription.kt index d7cd6b7a..87f78d58 100644 --- a/app/src/main/java/com/greenart7c3/nostrsigner/service/ProfileSubscription.kt +++ b/app/src/main/java/com/greenart7c3/nostrsigner/service/ProfileSubscription.kt @@ -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() + private val subIds = ConcurrentHashMap() // hexKey -> subId of the kind-10002 relay list subscription that runs before the metadata one - private val relayListSubIds = mutableMapOf() - private val relaysPerSubId = mutableMapOf>() - private val timeoutJobs = mutableMapOf() + private val relayListSubIds = ConcurrentHashMap() + private val relaysPerSubId = ConcurrentHashMap>() + private val timeoutJobs = ConcurrentHashMap() // 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() + private val accounts = ConcurrentHashMap() 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().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().apply { addAll(profileFilter.keys) } client.subscribe(subId, profileFilter) timeoutJobs[subId] = scope.launch { delay(EOSE_TIMEOUT_MS) diff --git a/app/src/test/java/com/greenart7c3/nostrsigner/service/NotificationSubscriptionTest.kt b/app/src/test/java/com/greenart7c3/nostrsigner/service/NotificationSubscriptionTest.kt new file mode 100644 index 00000000..d2ae22c5 --- /dev/null +++ b/app/src/test/java/com/greenart7c3/nostrsigner/service/NotificationSubscriptionTest.kt @@ -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()) + 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() + every { database.dao() } returns dao + + val amber = mockk(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() + val event = mockk(relaxed = true) + + val errors = Collections.synchronizedList(mutableListOf()) + val threads = mutableListOf() + + 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(), + ) + } +} diff --git a/app/src/test/java/com/greenart7c3/nostrsigner/service/ProfileSubscriptionTest.kt b/app/src/test/java/com/greenart7c3/nostrsigner/service/ProfileSubscriptionTest.kt new file mode 100644 index 00000000..2fcf9920 --- /dev/null +++ b/app/src/test/java/com/greenart7c3/nostrsigner/service/ProfileSubscriptionTest.kt @@ -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()) + 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(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()) + runBlocking { + repeat(200) { + val account = newTestAccount(scope) + tracked.add(account) + subscription.updateFilter(account) + } + } + + val errors = Collections.synchronizedList(mutableListOf()) + val threads = mutableListOf() + + 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() + every { relay.url } returns AmberSettings().defaultProfileRelays.first() + val event = mockk() + every { event.kind } returns -1 + + val accounts = (1..8).map { newTestAccount(scope) } + runBlocking { accounts.forEach { subscription.updateFilter(it) } } + + val errors = Collections.synchronizedList(mutableListOf()) + val threads = mutableListOf() + + 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(), + ) + } +} diff --git a/app/src/test/java/com/greenart7c3/nostrsigner/service/SubscriptionTestUtils.kt b/app/src/test/java/com/greenart7c3/nostrsigner/service/SubscriptionTestUtils.kt new file mode 100644 index 00000000..5148ed56 --- /dev/null +++ b/app/src/test/java/com/greenart7c3/nostrsigner/service/SubscriptionTestUtils.kt @@ -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(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, + ) +}