mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-05 19:28:25 +00:00
Code review: release awaiters + cap the NIP-44 negotiation wait
fix(nwc): release awaiters when the account scope is already dead fix(nwc): cap the NIP-44 negotiation wait and keep the fetch off the caller
This commit is contained in:
+37
-23
@@ -74,9 +74,7 @@ class NwcInfoCache(
|
||||
|
||||
private val cache = ConcurrentHashMap<HexKey, Entry>()
|
||||
|
||||
// Fetches in progress, keyed like [cache]. Every path that fetches goes through
|
||||
// [fetchOnce], so the startup warm-up, a payment waiting on a cold cache and the
|
||||
// notification watcher all join one request per wallet instead of racing.
|
||||
// Fetches in progress, keyed like [cache]. Every fetching path goes through [fetchOnce].
|
||||
private val inFlight = ConcurrentHashMap<HexKey, CompletableDeferred<NwcInfoEvent?>>()
|
||||
|
||||
private fun isFresh(entry: Entry): Boolean = now() - entry.fetchedAt < ttlSeconds
|
||||
@@ -92,7 +90,6 @@ class NwcInfoCache(
|
||||
fun refreshIfStale(uri: Nip47WalletConnect.Nip47URINorm) {
|
||||
val entry = cache[uri.pubKeyHex]
|
||||
if (entry != null && isFresh(entry)) return
|
||||
if (inFlight.containsKey(uri.pubKeyHex)) return
|
||||
|
||||
scope.launch(Dispatchers.IO) { fetchOnce(uri) }
|
||||
}
|
||||
@@ -101,6 +98,9 @@ class NwcInfoCache(
|
||||
* Suspends until a fresh-enough info event is available, fetching when the
|
||||
* entry is missing or expired. Returns the last cached (possibly stale) value
|
||||
* if the fetch fails.
|
||||
*
|
||||
* For a caller that must not act on a stale answer. No production caller needs
|
||||
* that today; prefer [currentOrFetch], which only waits on a cold cache.
|
||||
*/
|
||||
suspend fun getFresh(uri: Nip47WalletConnect.Nip47URINorm): NwcInfoEvent? {
|
||||
val entry = cache[uri.pubKeyHex]
|
||||
@@ -123,38 +123,52 @@ class NwcInfoCache(
|
||||
* background for next time.
|
||||
*/
|
||||
suspend fun currentOrFetch(uri: Nip47WalletConnect.Nip47URINorm): NwcInfoEvent? {
|
||||
val entry = cache[uri.pubKeyHex]
|
||||
if (entry != null) {
|
||||
refreshIfStale(uri)
|
||||
return entry.info
|
||||
}
|
||||
return fetchOnce(uri)
|
||||
val entry = cache[uri.pubKeyHex] ?: return fetchOnce(uri)
|
||||
refreshIfStale(uri)
|
||||
return entry.info
|
||||
}
|
||||
|
||||
/**
|
||||
* Runs [fetchAndStore] for [uri] exactly once, however many callers ask at
|
||||
* once; the rest await that one result.
|
||||
* once; every caller awaits that one result.
|
||||
*
|
||||
* Cancellation reaches only the caller that owns the fetch — awaiters are
|
||||
* separate coroutines that are still alive, and get the last cached value (or
|
||||
* null) rather than being cancelled along with it.
|
||||
* The fetch runs in this cache's own [scope], never in the caller's. A caller
|
||||
* on the payment path is a `viewModelScope` coroutine that dies when the user
|
||||
* backs out of the screen — owning the fetch there would abandon it, leave the
|
||||
* cache cold and make the next attempt pay the whole cost again. Here a caller
|
||||
* giving up cancels only its own `await`, and the fetch it started still lands.
|
||||
*/
|
||||
private suspend fun fetchOnce(uri: Nip47WalletConnect.Nip47URINorm): NwcInfoEvent? {
|
||||
val key = uri.pubKeyHex
|
||||
val ours = CompletableDeferred<NwcInfoEvent?>()
|
||||
// putIfAbsent returns the previous entry, so a non-null result means
|
||||
// someone else already owns this wallet's fetch.
|
||||
// someone else already started this wallet's fetch.
|
||||
inFlight.putIfAbsent(key, ours)?.let { return it.await() }
|
||||
|
||||
var info: NwcInfoEvent? = null
|
||||
try {
|
||||
info = fetchAndStore(uri)
|
||||
} finally {
|
||||
// Non-suspending, so awaiters are released even if we are cancelled.
|
||||
inFlight.remove(key, ours)
|
||||
ours.complete(info)
|
||||
val job =
|
||||
scope.launch(Dispatchers.IO) {
|
||||
var info: NwcInfoEvent? = null
|
||||
try {
|
||||
info = fetchAndStore(uri)
|
||||
} finally {
|
||||
// Non-suspending, so awaiters are released even if the fetch is cancelled.
|
||||
inFlight.remove(key, ours)
|
||||
ours.complete(info)
|
||||
}
|
||||
}
|
||||
|
||||
// A scope that was already cancelled — this account was logged off while a
|
||||
// payment was in flight — never runs the body, so the finally above never
|
||||
// releases anyone. Without this, the caller awaits forever and the abandoned
|
||||
// slot makes every later call for this wallet do the same. Fires immediately
|
||||
// when the job is already complete, and is a no-op on the normal path.
|
||||
job.invokeOnCompletion {
|
||||
if (!ours.isCompleted) {
|
||||
inFlight.remove(key, ours)
|
||||
ours.complete(null)
|
||||
}
|
||||
}
|
||||
return info
|
||||
return ours.await()
|
||||
}
|
||||
|
||||
private suspend fun fetchAndStore(uri: Nip47WalletConnect.Nip47URINorm): NwcInfoEvent? {
|
||||
|
||||
+19
-1
@@ -61,6 +61,7 @@ import kotlinx.coroutines.flow.flowOn
|
||||
import kotlinx.coroutines.flow.map
|
||||
import kotlinx.coroutines.flow.stateIn
|
||||
import kotlinx.coroutines.launch
|
||||
import kotlinx.coroutines.withTimeoutOrNull
|
||||
|
||||
/**
|
||||
* Manages NIP-47 (Nostr Wallet Connect) related signing operations and decryption cache for a given account.
|
||||
@@ -147,10 +148,20 @@ class NwcSignerState(
|
||||
* back to NIP-04 even against a wallet advertising `nip44_v2`. Only a cold
|
||||
* cache waits: a stale entry still says what the wallet advertises and is used
|
||||
* as-is while it refreshes in the background.
|
||||
*
|
||||
* The wait is capped, because this sits in front of a payment the user has
|
||||
* already tapped. Its own fetch is bounded only by the relay accessory's 30s
|
||||
* idle window, and the no-response timer below does not start until this
|
||||
* returns. On expiry we send NIP-04 for this one request rather than hold the
|
||||
* tap; the fetch keeps running in the cache's scope, so the next request gets
|
||||
* the negotiated scheme. Never bound the fetch itself instead — a null from it
|
||||
* is cached as a definitive "no info event" for the whole TTL, which would pin
|
||||
* the wallet to NIP-04 for days.
|
||||
*/
|
||||
private suspend fun prefersNip44(uri: Nip47WalletConnect.Nip47URINorm?): Boolean {
|
||||
uri ?: return false
|
||||
return infoCache?.currentOrFetch(uri)?.encryptionSchemes()?.any { it.equals("nip44_v2", ignoreCase = true) } ?: false
|
||||
val info = withTimeoutOrNull(NIP44_NEGOTIATION_WAIT_MS) { infoCache?.currentOrFetch(uri) }
|
||||
return info?.encryptionSchemes()?.any { it.equals("nip44_v2", ignoreCase = true) } == true
|
||||
}
|
||||
|
||||
fun hasWalletConnectSetup(): Boolean = settings.nwcWallets.value.isNotEmpty()
|
||||
@@ -336,5 +347,12 @@ class NwcSignerState(
|
||||
const val NWC_RESPONSE_TIMEOUT_SECONDS = 60
|
||||
|
||||
const val NWC_RESPONSE_TIMEOUT_MS = NWC_RESPONSE_TIMEOUT_SECONDS * 1000L
|
||||
|
||||
/**
|
||||
* How long a request will wait for a cold info cache before falling back to
|
||||
* NIP-04. Comfortably over a healthy single-relay round trip, far under the
|
||||
* 30s the fetch itself would otherwise allow in front of a payment tap.
|
||||
*/
|
||||
const val NIP44_NEGOTIATION_WAIT_MS = 3_000L
|
||||
}
|
||||
}
|
||||
|
||||
+58
-11
@@ -28,8 +28,13 @@ import kotlinx.coroutines.CoroutineScope
|
||||
import kotlinx.coroutines.Dispatchers
|
||||
import kotlinx.coroutines.async
|
||||
import kotlinx.coroutines.awaitAll
|
||||
import kotlinx.coroutines.cancel
|
||||
import kotlinx.coroutines.cancelAndJoin
|
||||
import kotlinx.coroutines.coroutineScope
|
||||
import kotlinx.coroutines.delay
|
||||
import kotlinx.coroutines.launch
|
||||
import kotlinx.coroutines.runBlocking
|
||||
import kotlinx.coroutines.withTimeout
|
||||
import org.junit.Assert.assertEquals
|
||||
import org.junit.Assert.assertNotNull
|
||||
import org.junit.Assert.assertNull
|
||||
@@ -172,7 +177,7 @@ class NwcInfoCacheTest {
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `concurrent getFresh callers share one fetch`() =
|
||||
fun concurrentGetFreshCallersShareOneFetch() =
|
||||
runBlocking {
|
||||
val fetcher = GatedFetch(info("pay_invoice"))
|
||||
val cache = NwcInfoCache(fetch = fetcher::fetch, scope = scope, ttlSeconds = 100, now = { clock })
|
||||
@@ -191,7 +196,7 @@ class NwcInfoCacheTest {
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `getFresh joins a fetch already started by refreshIfStale`() =
|
||||
fun getFreshJoinsFetchAlreadyStartedByRefreshIfStale() =
|
||||
runBlocking {
|
||||
val fetcher = GatedFetch(info("pay_invoice"))
|
||||
val cache = NwcInfoCache(fetch = fetcher::fetch, scope = scope, ttlSeconds = 100, now = { clock })
|
||||
@@ -213,7 +218,7 @@ class NwcInfoCacheTest {
|
||||
// --- currentOrFetch: wait only when there is nothing cached at all ---
|
||||
|
||||
@Test
|
||||
fun `currentOrFetch waits for the first fetch when nothing is cached`() =
|
||||
fun currentOrFetchWaitsWhenNothingIsCached() =
|
||||
runBlocking {
|
||||
val fetcher = GatedFetch(info("pay_invoice"))
|
||||
val cache = NwcInfoCache(fetch = fetcher::fetch, scope = scope, ttlSeconds = 100, now = { clock })
|
||||
@@ -231,13 +236,16 @@ class NwcInfoCacheTest {
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `currentOrFetch returns a stale entry without waiting`() =
|
||||
fun currentOrFetchReturnsStaleEntryWithoutWaiting() =
|
||||
runBlocking {
|
||||
var calls = 0
|
||||
// The background refresh is made to hang, so waiting on it would hang the
|
||||
// test. Returning at all is the assertion.
|
||||
val hang = CompletableDeferred<Unit>()
|
||||
val calls = AtomicInteger(0)
|
||||
val cache =
|
||||
NwcInfoCache(
|
||||
fetch = {
|
||||
calls++
|
||||
if (calls.incrementAndGet() > 1) hang.await()
|
||||
info("pay_invoice")
|
||||
},
|
||||
scope = scope,
|
||||
@@ -246,12 +254,51 @@ class NwcInfoCacheTest {
|
||||
)
|
||||
|
||||
assertNotNull(cache.currentOrFetch(uri("wallet1")))
|
||||
assertEquals(1, calls)
|
||||
assertEquals(1, calls.get())
|
||||
|
||||
// Past the TTL the entry is stale, but it still answers the question it
|
||||
// is asked (which encryption the wallet advertises), so the caller must
|
||||
// get it back without a round-trip. The background refresh self-heals.
|
||||
// Past the TTL the entry is stale, but it still says which encryption the
|
||||
// wallet advertises, so it comes back as-is while the refresh runs on.
|
||||
clock += 1_000
|
||||
assertNotNull(cache.currentOrFetch(uri("wallet1")))
|
||||
val stale = cache.currentOrFetch(uri("wallet1"))
|
||||
hang.complete(Unit)
|
||||
assertNotNull(stale)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun cancellingTheCallerDoesNotAbortTheSharedFetch() =
|
||||
runBlocking {
|
||||
// A payment coroutine that asks first must not own the fetch: backing out
|
||||
// of the payment screen would cancel it, leaving the cache cold so the
|
||||
// next attempt pays the whole cost again.
|
||||
val fetcher = GatedFetch(info("pay_invoice"))
|
||||
val cache = NwcInfoCache(fetch = fetcher::fetch, scope = scope, ttlSeconds = 100, now = { clock })
|
||||
|
||||
val caller = launch(Dispatchers.IO) { cache.currentOrFetch(uri("wallet1")) }
|
||||
fetcher.started.await()
|
||||
caller.cancelAndJoin()
|
||||
fetcher.release.complete(Unit)
|
||||
|
||||
withTimeout(5_000) {
|
||||
while (cache.current(uri("wallet1")) == null) delay(5)
|
||||
}
|
||||
assertEquals("the abandoned fetch still completed and warmed the cache", 1, fetcher.calls.get())
|
||||
}
|
||||
|
||||
@Test
|
||||
fun aDeadScopeReleasesCallersInsteadOfHanging() =
|
||||
runBlocking {
|
||||
// Account.scope is cancelled on logout (AccountCacheState.removeAccount),
|
||||
// which runs asynchronously and can race a payment already in flight.
|
||||
// Launching into it does not run the body, so nothing would complete the
|
||||
// deferred a caller is awaiting.
|
||||
val dead = CoroutineScope(Dispatchers.IO).also { it.cancel() }
|
||||
val fetcher = GatedFetch(info("pay_invoice"))
|
||||
val cache = NwcInfoCache(fetch = fetcher::fetch, scope = dead, ttlSeconds = 100, now = { clock })
|
||||
|
||||
assertNull(withTimeout(2_000) { cache.currentOrFetch(uri("wallet1")) })
|
||||
assertEquals(0, fetcher.calls.get())
|
||||
|
||||
// and the abandoned slot must not poison every later call for that wallet
|
||||
assertNull(withTimeout(2_000) { cache.currentOrFetch(uri("wallet1")) })
|
||||
}
|
||||
}
|
||||
|
||||
+1
-1
@@ -49,7 +49,7 @@ class OkHttpLnurlEndpointResolver(
|
||||
) : LnurlEndpointResolver {
|
||||
private val mapper = jacksonObjectMapper()
|
||||
|
||||
override suspend fun resolve(lnurlpUrl: String): LnurlEndpointInfo? = LnurlEndpointCache.getOrFetch(lnurlpUrl) { fetch(it) }
|
||||
override suspend fun resolve(lnurlpUrl: String): LnurlEndpointInfo? = LnurlEndpointCache.getOrFetch(lnurlpUrl, ::fetch)
|
||||
|
||||
private suspend fun fetch(url: String): LnurlEndpointInfo? =
|
||||
withContext(Dispatchers.IO) {
|
||||
|
||||
+18
-1
@@ -97,7 +97,24 @@ class OkHttpLnurlEndpointResolverTest {
|
||||
@Test
|
||||
fun `a burst of callers makes one http request`() =
|
||||
runBlocking {
|
||||
val results = burst(20)
|
||||
// The gate matters: the claim under test is that these callers were all
|
||||
// in flight together, and a straggler arriving after the first response
|
||||
// landed would legitimately issue a second request.
|
||||
val client = OkHttpClient.Builder().addInterceptor(interceptor).build()
|
||||
val resolver = OkHttpLnurlEndpointResolver { client }
|
||||
val gate = CompletableDeferred<Unit>()
|
||||
val results =
|
||||
coroutineScope {
|
||||
val callers =
|
||||
(0 until 20).map {
|
||||
async(Dispatchers.IO) {
|
||||
gate.await()
|
||||
resolver.resolve(URL)
|
||||
}
|
||||
}
|
||||
gate.complete(Unit)
|
||||
callers.awaitAll()
|
||||
}
|
||||
|
||||
assertEquals(1, interceptor.calls.get(), "20 concurrent callers must share one request")
|
||||
assertNotNull(results.first(), "the shared fetch resolved")
|
||||
|
||||
+4
-4
@@ -100,10 +100,10 @@ object LnurlEndpointCache {
|
||||
try {
|
||||
info = fetch(url)?.also { cache.put(key, it) }
|
||||
} finally {
|
||||
// Populate the cache before releasing the flight, so a caller arriving
|
||||
// in between reads the result instead of starting a second fetch.
|
||||
// Both calls are non-suspending, so this still runs — and still
|
||||
// unblocks the awaiters — if we are cancelled mid-fetch.
|
||||
// The cache is populated before the flight is released, so a caller
|
||||
// arriving in between reads the result instead of starting a second
|
||||
// fetch. Non-suspending, so awaiters are released even if the fetch is
|
||||
// cancelled.
|
||||
inFlight.remove(key, ours)
|
||||
ours.complete(info)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user