diff --git a/amethyst/src/main/java/com/vitorpamplona/amethyst/model/nip47WalletConnect/NwcInfoCache.kt b/amethyst/src/main/java/com/vitorpamplona/amethyst/model/nip47WalletConnect/NwcInfoCache.kt index 09c4d10423..603f5b6bd6 100644 --- a/amethyst/src/main/java/com/vitorpamplona/amethyst/model/nip47WalletConnect/NwcInfoCache.kt +++ b/amethyst/src/main/java/com/vitorpamplona/amethyst/model/nip47WalletConnect/NwcInfoCache.kt @@ -74,9 +74,7 @@ class NwcInfoCache( private val cache = ConcurrentHashMap() - // 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>() 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() // 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? { diff --git a/amethyst/src/main/java/com/vitorpamplona/amethyst/model/nip47WalletConnect/NwcSignerState.kt b/amethyst/src/main/java/com/vitorpamplona/amethyst/model/nip47WalletConnect/NwcSignerState.kt index c6e9af96cf..ca99213d1f 100644 --- a/amethyst/src/main/java/com/vitorpamplona/amethyst/model/nip47WalletConnect/NwcSignerState.kt +++ b/amethyst/src/main/java/com/vitorpamplona/amethyst/model/nip47WalletConnect/NwcSignerState.kt @@ -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 } } diff --git a/amethyst/src/test/java/com/vitorpamplona/amethyst/model/nip47WalletConnect/NwcInfoCacheTest.kt b/amethyst/src/test/java/com/vitorpamplona/amethyst/model/nip47WalletConnect/NwcInfoCacheTest.kt index 8a6a7f45a4..584df37633 100644 --- a/amethyst/src/test/java/com/vitorpamplona/amethyst/model/nip47WalletConnect/NwcInfoCacheTest.kt +++ b/amethyst/src/test/java/com/vitorpamplona/amethyst/model/nip47WalletConnect/NwcInfoCacheTest.kt @@ -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() + 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")) }) } } diff --git a/commons/src/jvmAndroid/kotlin/com/vitorpamplona/amethyst/commons/service/lnurl/OkHttpLnurlEndpointResolver.kt b/commons/src/jvmAndroid/kotlin/com/vitorpamplona/amethyst/commons/service/lnurl/OkHttpLnurlEndpointResolver.kt index a9802aec99..2fbdf41481 100644 --- a/commons/src/jvmAndroid/kotlin/com/vitorpamplona/amethyst/commons/service/lnurl/OkHttpLnurlEndpointResolver.kt +++ b/commons/src/jvmAndroid/kotlin/com/vitorpamplona/amethyst/commons/service/lnurl/OkHttpLnurlEndpointResolver.kt @@ -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) { diff --git a/commons/src/jvmTest/kotlin/com/vitorpamplona/amethyst/commons/service/lnurl/OkHttpLnurlEndpointResolverTest.kt b/commons/src/jvmTest/kotlin/com/vitorpamplona/amethyst/commons/service/lnurl/OkHttpLnurlEndpointResolverTest.kt index 26aa7fc463..7e14f1df49 100644 --- a/commons/src/jvmTest/kotlin/com/vitorpamplona/amethyst/commons/service/lnurl/OkHttpLnurlEndpointResolverTest.kt +++ b/commons/src/jvmTest/kotlin/com/vitorpamplona/amethyst/commons/service/lnurl/OkHttpLnurlEndpointResolverTest.kt @@ -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() + 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") diff --git a/quartz/src/jvmAndroid/kotlin/com/vitorpamplona/quartz/nip57Zaps/validate/LnurlEndpointCache.kt b/quartz/src/jvmAndroid/kotlin/com/vitorpamplona/quartz/nip57Zaps/validate/LnurlEndpointCache.kt index 5a58ed2410..c7fe418607 100644 --- a/quartz/src/jvmAndroid/kotlin/com/vitorpamplona/quartz/nip57Zaps/validate/LnurlEndpointCache.kt +++ b/quartz/src/jvmAndroid/kotlin/com/vitorpamplona/quartz/nip57Zaps/validate/LnurlEndpointCache.kt @@ -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) }