From e2743ed0b63fa3f2e0c0b7044a5c13324807ac01 Mon Sep 17 00:00:00 2001 From: Claude Date: Fri, 17 Jul 2026 02:10:18 +0000 Subject: [PATCH] fix: address pre-merge audit findings for geohash chat Correctness: - GeoRelayDirectory.relays is now @Volatile; on the process-wide `shared` directory the CSV refresh was written by one thread and read by others with no memory barrier, so readers could keep using the FALLBACK list forever and never route to the correct rendezvous relays. - sendPostSync bails before cancel() when a geohash cell has no resolvable relays, so the composer text + draft are preserved instead of the message being silently dropped with its draft deleted. - Teleport detection compares on the common geohash prefix; a cell finer than the fixed 8-char device fix could never be a startsWith prefix, so the user was wrongly marked teleported even when physically present. Performance: - GeoRelayDirectory.closest precomputes each relay's great-circle distance once instead of recomputing the trig inside the sort comparator (was O(n log n) haversine calls over the ~370-relay directory). - GeohashChatChannel.relays() memoizes the derived set, invalidated by a new directory version token, instead of re-sorting the whole directory (and allocating a fresh Set) on every call. - filterFollowingGeohashChats groups cells by relay into one filter each (g = [cells]) rather than one REQ per (cell, relay). Leak/thread-safety: - FollowingGeohashChatSubAssembler.userJobMap is a ConcurrentHashMap and endSub now removes the entry (it previously cancelled the jobs but left the stale entry behind). Co-Authored-By: Claude Opus 4.8 Claude-Session: https://claude.ai/code/session_0172JoMccseEKenyWan6txWV --- .../chats/geohashChat/GeohashChatScreen.kt | 8 ++++- .../send/ChannelNewMessageViewModel.kt | 4 +++ .../datasource/FilterFollowingGeohashChats.kt | 32 ++++++++++++------- .../FollowingGeohashChatSubAssembler.kt | 10 +++--- .../model/geohashChat/GeohashChatChannel.kt | 18 ++++++++++- .../service/georelay/GeoRelayDirectory.kt | 22 ++++++++++--- 6 files changed, 72 insertions(+), 22 deletions(-) diff --git a/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/chats/geohashChat/GeohashChatScreen.kt b/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/chats/geohashChat/GeohashChatScreen.kt index 54f7c60898..4c2ee851b5 100644 --- a/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/chats/geohashChat/GeohashChatScreen.kt +++ b/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/chats/geohashChat/GeohashChatScreen.kt @@ -168,7 +168,13 @@ private fun GeohashChatRoom( val isTeleported = remember(deviceLocation, geohash, teleported) { when (val loc = deviceLocation) { - is LocationState.LocationResult.Success -> !loc.geoHash.toString().startsWith(geohash) + is LocationState.LocationResult.Success -> { + // Compare on the common prefix: the device fix is fixed-precision (8 chars), so a + // finer (longer) cell can never be a prefix of it — treat a shared prefix as "here". + val device = loc.geoHash.toString() + val prefix = minOf(device.length, geohash.length) + device.take(prefix) != geohash.take(prefix) + } else -> teleported } } diff --git a/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/chats/publicChannels/send/ChannelNewMessageViewModel.kt b/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/chats/publicChannels/send/ChannelNewMessageViewModel.kt index 2c79c5cb68..17a13aee09 100644 --- a/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/chats/publicChannels/send/ChannelNewMessageViewModel.kt +++ b/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/chats/publicChannels/send/ChannelNewMessageViewModel.kt @@ -339,6 +339,10 @@ open class ChannelNewMessageViewModel : val template = createTemplate() ?: return val channelRelays = channel?.relays() ?: emptySet() + // A geohash cell with no resolvable relays has nowhere to publish. Bail before cancel() clears + // the composer, so the user keeps their text (and draft) to retry rather than losing it silently. + if (channel is GeohashChatChannel && channelRelays.isEmpty()) return + val draftToDelete = draftNote cancel() diff --git a/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/chats/rooms/datasource/FilterFollowingGeohashChats.kt b/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/chats/rooms/datasource/FilterFollowingGeohashChats.kt index 60e3814a2b..bb63652938 100644 --- a/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/chats/rooms/datasource/FilterFollowingGeohashChats.kt +++ b/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/chats/rooms/datasource/FilterFollowingGeohashChats.kt @@ -25,6 +25,7 @@ import com.vitorpamplona.amethyst.service.relays.SincePerRelayMap import com.vitorpamplona.quartz.experimental.bitchat.geohash.GeohashChatEvent import com.vitorpamplona.quartz.nip01Core.relay.client.pool.RelayBasedFilter import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter +import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl /** * REQ for the user's joined geohash location channels: for each cell, subscribe @@ -39,18 +40,25 @@ fun filterFollowingGeohashChats( ): List? { if (geohashes.isEmpty()) return null - return geohashes.flatMap { geohash -> - GeohashRelays.closestRelays(geohash).map { relay -> - RelayBasedFilter( - relay = relay, - filter = - Filter( - kinds = listOf(GeohashChatEvent.KIND), - tags = mapOf("g" to listOf(geohash)), - limit = 100, - since = since?.get(relay)?.time, - ), - ) + // Group cells by their nearest relays, so cells that share a relay collapse into a single filter + // (g = [all those cells]) instead of one REQ per (cell, relay). + val cellsByRelay = LinkedHashMap>() + geohashes.forEach { geohash -> + GeohashRelays.closestRelays(geohash).forEach { relay -> + cellsByRelay.getOrPut(relay) { mutableListOf() }.add(geohash) } } + + return cellsByRelay.map { (relay, cells) -> + RelayBasedFilter( + relay = relay, + filter = + Filter( + kinds = listOf(GeohashChatEvent.KIND), + tags = mapOf("g" to cells), + limit = 100, + since = since?.get(relay)?.time, + ), + ) + } } diff --git a/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/chats/rooms/datasource/FollowingGeohashChatSubAssembler.kt b/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/chats/rooms/datasource/FollowingGeohashChatSubAssembler.kt index 4438d0afb0..9186e0ccb3 100644 --- a/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/chats/rooms/datasource/FollowingGeohashChatSubAssembler.kt +++ b/amethyst/src/main/java/com/vitorpamplona/amethyst/ui/screen/loggedIn/chats/rooms/datasource/FollowingGeohashChatSubAssembler.kt @@ -27,6 +27,7 @@ import com.vitorpamplona.amethyst.service.relays.SincePerRelayMap import com.vitorpamplona.quartz.nip01Core.relay.client.INostrClient import com.vitorpamplona.quartz.nip01Core.relay.client.pool.RelayBasedFilter import com.vitorpamplona.quartz.nip01Core.relay.client.subscriptions.Subscription +import java.util.concurrent.ConcurrentHashMap import kotlinx.coroutines.Dispatchers import kotlinx.coroutines.FlowPreview import kotlinx.coroutines.Job @@ -55,12 +56,13 @@ class FollowingGeohashChatSubAssembler( override fun user(key: ChatroomListState) = key.account.userProfile() - val userJobMap = mutableMapOf>() + private val userJobMap = ConcurrentHashMap>() @OptIn(FlowPreview::class) override fun newSub(key: ChatroomListState): Subscription { - userJobMap[key.account.userProfile()]?.forEach { it.cancel() } - userJobMap[key.account.userProfile()] = + val user = key.account.userProfile() + userJobMap.remove(user)?.forEach { it.cancel() } + userJobMap[user] = listOf( // Rebuild when the joined set changes. key.account.scope.launch(Dispatchers.IO) { @@ -82,6 +84,6 @@ class FollowingGeohashChatSubAssembler( subId: String, ) { super.endSub(key, subId) - userJobMap[key]?.forEach { it.cancel() } + userJobMap.remove(key)?.forEach { it.cancel() } } } diff --git a/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/model/geohashChat/GeohashChatChannel.kt b/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/model/geohashChat/GeohashChatChannel.kt index 8f14cfdcba..0cbf00b27a 100644 --- a/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/model/geohashChat/GeohashChatChannel.kt +++ b/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/model/geohashChat/GeohashChatChannel.kt @@ -28,6 +28,7 @@ import com.vitorpamplona.amethyst.commons.util.KmpLock import com.vitorpamplona.amethyst.commons.util.withLock import com.vitorpamplona.quartz.nip01Core.core.HexKey import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl +import kotlin.concurrent.Volatile /** * A public geohash location chat channel (Bitchat-interoperable), keyed by the @@ -49,7 +50,22 @@ class GeohashChatChannel( * cell has a well-defined relay set derived from its coordinates, so the * subscription layer can reach it even before the first message arrives. */ - override fun relays(): Set = GeoRelayDirectory.shared.closestRelays(geohash).toSet() + @Volatile private var cachedRelays: Set? = null + + @Volatile private var cachedRelaysVersion: Int = -1 + + override fun relays(): Set { + val dir = GeoRelayDirectory.shared + val v = dir.version + // The closest-relays computation is a full sort of the ~370-relay directory; memoize it and + // recompute only when the directory changes (its version token moves). Two threads racing on a + // cold cache just both compute the same deterministic result, so no lock is needed. + cachedRelays?.let { if (cachedRelaysVersion == v) return it } + val computed = dir.closestRelays(geohash).toSet() + cachedRelays = computed + cachedRelaysVersion = v + return computed + } fun anyNameStartsWith(prefix: String): Boolean = geohash.contains(prefix, true) diff --git a/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/service/georelay/GeoRelayDirectory.kt b/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/service/georelay/GeoRelayDirectory.kt index 219da21358..1707dfe243 100644 --- a/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/service/georelay/GeoRelayDirectory.kt +++ b/commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/service/georelay/GeoRelayDirectory.kt @@ -23,6 +23,7 @@ package com.vitorpamplona.amethyst.commons.service.georelay import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl import com.vitorpamplona.quartz.nip01Core.relay.normalizer.RelayUrlNormalizer import com.vitorpamplona.quartz.nip01Core.tags.geohash.GeoHash +import kotlin.concurrent.Volatile import kotlin.math.asin import kotlin.math.cos import kotlin.math.min @@ -48,11 +49,21 @@ import kotlin.math.sqrt class GeoRelayDirectory( initial: List = FALLBACK, ) { - private var relays: List = initial + // Read from subscription/UI threads, replaced by the CSV-loader coroutine — @Volatile so the + // refresh is actually visible to readers (otherwise they can keep seeing FALLBACK forever). + @Volatile private var relays: List = initial + + // Bumped whenever the directory contents change, so cheap consumers (e.g. GeohashChatChannel) can + // memoize a derived relay set and recompute only when this token moves. + @Volatile var version: Int = 0 + private set /** Replace the directory contents (e.g. after a successful CSV refresh). No-op on empty input. */ fun setRelays(list: List) { - if (list.isNotEmpty()) relays = list + if (list.isNotEmpty()) { + relays = list + version++ + } } fun snapshot(): List = relays @@ -90,8 +101,11 @@ class GeoRelayDirectory( count: Int = DEFAULT_COUNT, ): List = relays - .sortedWith(compareBy({ haversineKm(lat, lon, it.latitude, it.longitude) }, { it.host })) - .map { it.relay } + // Compute the great-circle distance once per relay; folding it into the comparator + // selector instead would re-run the trig on every comparison (O(n log n) times). + .map { it to haversineKm(lat, lon, it.latitude, it.longitude) } + .sortedWith(compareBy({ it.second }, { it.first.host })) + .map { it.first.relay } // The directory lists some hosts both bare and with an explicit :443 (the // wss default) — those collapse to one relay, so drop the duplicate before // taking count and never spend two of the N slots on the same endpoint.