mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-06 03:38:23 +00:00
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 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_0172JoMccseEKenyWan6txWV
This commit is contained in:
+7
-1
@@ -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
|
||||
}
|
||||
}
|
||||
|
||||
+4
@@ -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()
|
||||
|
||||
|
||||
+20
-12
@@ -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<RelayBasedFilter>? {
|
||||
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<NormalizedRelayUrl, MutableList<String>>()
|
||||
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,
|
||||
),
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
+6
-4
@@ -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<User, List<Job>>()
|
||||
private val userJobMap = ConcurrentHashMap<User, List<Job>>()
|
||||
|
||||
@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() }
|
||||
}
|
||||
}
|
||||
|
||||
+17
-1
@@ -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<NormalizedRelayUrl> = GeoRelayDirectory.shared.closestRelays(geohash).toSet()
|
||||
@Volatile private var cachedRelays: Set<NormalizedRelayUrl>? = null
|
||||
|
||||
@Volatile private var cachedRelaysVersion: Int = -1
|
||||
|
||||
override fun relays(): Set<NormalizedRelayUrl> {
|
||||
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)
|
||||
|
||||
|
||||
+18
-4
@@ -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<GeoRelay> = FALLBACK,
|
||||
) {
|
||||
private var relays: List<GeoRelay> = 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<GeoRelay> = 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<GeoRelay>) {
|
||||
if (list.isNotEmpty()) relays = list
|
||||
if (list.isNotEmpty()) {
|
||||
relays = list
|
||||
version++
|
||||
}
|
||||
}
|
||||
|
||||
fun snapshot(): List<GeoRelay> = relays
|
||||
@@ -90,8 +101,11 @@ class GeoRelayDirectory(
|
||||
count: Int = DEFAULT_COUNT,
|
||||
): List<NormalizedRelayUrl> =
|
||||
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.
|
||||
|
||||
Reference in New Issue
Block a user