refactor(zaps): move LNURL fetch dedup onto LnurlEndpointCache

Single-flight landed inside OkHttpLnurlEndpointResolver, which put the two
halves of one mechanism — "resolve this URL exactly once" — in two modules.
The flight map had to call LnurlForm.normalizeUrl purely to match a keying
detail private to LnurlEndpointCache in quartz. Nothing documented or
enforced that: if the cache changed its canonicalisation, the map would
silently stop deduplicating and no test would fail.

Move it onto the cache as getOrFetch(url, fetch). The key is now computed
once and shared by the lookup, the flight map and the store, so they cannot
disagree. Dedup also becomes process-wide, matching the resource it
protects — a stranger's /.well-known/ endpoint — rather than being scoped to
one resolver instance; clear() resets both maps. The resolver drops to a
one-line delegation and keeps only the HTTP half. Same shape as NwcInfoCache,
which already pairs a cache with an in-flight map and an injected fetch.

Mechanism tests move to quartz beside the cache, using delay() rather than
a blocking sleep. The commons test keeps the one claim it uniquely makes:
that the resolver really routes through the cache over a real OkHttp client.

No behaviour change. Verified by mutation: removing single-flight, keying
the flight map on the raw URL, never releasing the slot, and making the
resolver bypass the cache each fail exactly the test that covers them.
This commit is contained in:
davotoula
2026-08-28 10:37:19 +02:00
parent 319348de9a
commit b29519e4e6
4 changed files with 204 additions and 125 deletions
@@ -24,14 +24,11 @@ import com.fasterxml.jackson.module.kotlin.jacksonObjectMapper
import com.vitorpamplona.quartz.nip57Zaps.validate.LnurlEndpointCache
import com.vitorpamplona.quartz.nip57Zaps.validate.LnurlEndpointInfo
import com.vitorpamplona.quartz.nip57Zaps.validate.LnurlEndpointResolver
import com.vitorpamplona.quartz.nip57Zaps.validate.LnurlForm
import kotlinx.coroutines.CompletableDeferred
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.withContext
import okhttp3.OkHttpClient
import okhttp3.Request
import okhttp3.coroutines.executeAsync
import java.util.concurrent.ConcurrentHashMap
import kotlin.coroutines.cancellation.CancellationException
/**
@@ -44,45 +41,15 @@ import kotlin.coroutines.cancellation.CancellationException
* caches the result, and returns it. Returns null on HTTP / parse failure;
* callers should treat that as "validation unavailable" rather than "invalid".
*
* Fetches are single-flighted per URL. A feed carrying twenty receipts for one
* lightning address hands us twenty concurrent [resolve] calls before the first
* fetch can populate the cache; without this, that is twenty requests to a
* stranger's unthrottled `/.well-known/` endpoint for one user action.
* Caching and fetch deduplication both belong to [LnurlEndpointCache]; this
* class is only the HTTP half.
*/
class OkHttpLnurlEndpointResolver(
private val okHttpClient: (String) -> OkHttpClient,
) : LnurlEndpointResolver {
private val mapper = jacksonObjectMapper()
/**
* Fetches in progress, keyed the same way [LnurlEndpointCache] keys itself so
* two spellings of one address share a flight. An entry lives only for the
* duration of its fetch: the winner removes it in a `finally`, so a failed
* fetch is retried by the next caller rather than being remembered as null.
*/
private val inFlight = ConcurrentHashMap<String, CompletableDeferred<LnurlEndpointInfo?>>()
override suspend fun resolve(lnurlpUrl: String): LnurlEndpointInfo? {
LnurlEndpointCache.get(lnurlpUrl)?.let { return it }
val key = LnurlForm.normalizeUrl(lnurlpUrl)
val ours = CompletableDeferred<LnurlEndpointInfo?>()
// Whoever wins putIfAbsent owns the fetch; everyone else awaits its result.
inFlight.putIfAbsent(key, ours)?.let { return it.await() }
var info: LnurlEndpointInfo? = null
try {
info = fetch(lnurlpUrl)
if (info != null) LnurlEndpointCache.put(lnurlpUrl, info)
} finally {
// Cache first, then release, so a caller arriving in between reads the
// cache rather than starting a second flight. Non-suspending, so it
// still runs — and still unblocks the awaiters — if we are cancelled.
inFlight.remove(key, ours)
ours.complete(info)
}
return info
}
override suspend fun resolve(lnurlpUrl: String): LnurlEndpointInfo? = LnurlEndpointCache.getOrFetch(lnurlpUrl) { fetch(it) }
private suspend fun fetch(url: String): LnurlEndpointInfo? =
withContext(Dispatchers.IO) {
@@ -42,131 +42,73 @@ import kotlin.test.assertNotNull
import kotlin.test.assertTrue
/**
* A zap-receipt burst hands the resolver N concurrent calls for the same
* lightning address. Each one is meant to await a single fetch, not start its
* own — the endpoint the fetch hits belongs to somebody else's server.
* Deduplication itself belongs to `LnurlEndpointCache` and is pinned by
* `LnurlEndpointCacheTest` in quartz. What is left to prove here is that this
* resolver actually routes through it, so a zap-receipt burst turns into one
* request on a real OkHttp client rather than one per receipt.
*/
class OkHttpLnurlEndpointResolverTest {
private val url = "https://example.com/.well-known/lnurlp/vitor"
/** Counts requests and holds each open long enough for the burst to pile up behind it. */
private class CountingInterceptor : Interceptor {
val calls = AtomicInteger(0)
/** Counts fetches and holds each one open long enough for the burst to pile up behind it. */
private class CountingInterceptor(
val calls: AtomicInteger = AtomicInteger(0),
private val body: String = """{"nostrPubkey":"aabbcc","allowsNostr":true}""",
private val delayMs: Long = 200,
private val failFirst: Int = 0,
) : Interceptor {
override fun intercept(chain: Interceptor.Chain): Response {
val n = calls.incrementAndGet()
Thread.sleep(delayMs)
if (n <= failFirst) {
return Response
.Builder()
.request(chain.request())
.protocol(Protocol.HTTP_1_1)
.code(503)
.message("Service Unavailable")
.body("".toResponseBody("text/plain".toMediaType()))
.build()
}
calls.incrementAndGet()
Thread.sleep(HOLD_MS)
return Response
.Builder()
.request(chain.request())
.protocol(Protocol.HTTP_1_1)
.code(200)
.message("OK")
.body(body.toResponseBody("application/json".toMediaType()))
.body(BODY.toResponseBody("application/json".toMediaType()))
.build()
}
}
private fun resolverWith(interceptor: Interceptor): OkHttpLnurlEndpointResolver {
val client = OkHttpClient.Builder().addInterceptor(interceptor).build()
return OkHttpLnurlEndpointResolver { client }
}
/**
* Releases [n] callers at once and collects what they got.
*
* The gate matters: the assertion under test is "these callers were all in
* flight together", and without it a straggler that arrives after the
* winner's fetch has already returned legitimately starts a second one. That
* would be a slow test machine failing the test, not a regression.
* Releases [n] callers at once. 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.
*/
private suspend fun burst(
n: Int,
call: suspend (Int) -> LnurlEndpointInfo?,
): List<LnurlEndpointInfo?> =
private suspend fun burst(n: Int): List<LnurlEndpointInfo?> =
coroutineScope {
val client = OkHttpClient.Builder().addInterceptor(interceptor).build()
val resolver = OkHttpLnurlEndpointResolver { client }
val gate = CompletableDeferred<Unit>()
val callers =
(0 until n).map { i ->
(0 until n).map {
async(Dispatchers.IO) {
gate.await()
call(i)
resolver.resolve(URL)
}
}
gate.complete(Unit)
callers.awaitAll()
}
private val interceptor = CountingInterceptor()
@BeforeTest
fun clearCache() {
LnurlEndpointCache.clear()
}
@Test
fun concurrentCallersForTheSameUrlShareOneFetch() =
fun `a burst of callers makes one http request`() =
runBlocking {
val interceptor = CountingInterceptor()
val resolver = resolverWith(interceptor)
val results = burst(20)
val results = burst(20) { resolver.resolve(url) }
assertEquals(1, interceptor.calls.get(), "20 concurrent callers must share one fetch")
assertEquals(20, results.size)
assertTrue(results.all { it != null && it == results.first() }, "every caller gets the same result")
assertEquals(1, interceptor.calls.get(), "20 concurrent callers must share one request")
assertNotNull(results.first(), "the shared fetch resolved")
assertTrue(results.all { it == results.first() }, "every caller gets the same result")
}
@Test
fun aFailedFetchIsNotRememberedAndTheNextCallerRetries() =
runBlocking {
// The server is down for the first burst. Nothing may be cached from
// that — otherwise one bad minute poisons the address until eviction.
val interceptor = CountingInterceptor(failFirst = 1)
val resolver = resolverWith(interceptor)
private companion object {
const val URL = "https://example.com/.well-known/lnurlp/vitor"
const val BODY = """{"nostrPubkey":"aabbcc","allowsNostr":true}"""
val firstBurst = burst(5) { resolver.resolve(url) }
assertEquals(1, interceptor.calls.get(), "the failing burst is one fetch, not five")
assertTrue(firstBurst.all { it == null }, "a 503 resolves to null, not a cached verdict")
val retry = resolver.resolve(url)
assertEquals(2, interceptor.calls.get(), "the next caller retries a failed address")
assertNotNull(retry, "the retry sees the recovered server")
resolver.resolve(url)
assertEquals(2, interceptor.calls.get(), "and the success is cached from then on")
}
@Test
fun twoSpellingsOfOneAddressShareOneFetch() =
runBlocking {
// LnurlEndpointCache keys case-insensitively on host and ignores a
// trailing slash. The in-flight map has to agree, or a burst that
// spells the host two ways is two stampedes instead of one.
val interceptor = CountingInterceptor(delayMs = 200)
val resolver = resolverWith(interceptor)
val spellings =
listOf(
"https://example.com/.well-known/lnurlp/vitor",
"https://EXAMPLE.com/.well-known/lnurlp/vitor",
"https://example.com/.well-known/lnurlp/vitor/",
)
val results = burst(12) { i -> resolver.resolve(spellings[i % 3]) }
assertEquals(1, interceptor.calls.get(), "host case and a trailing slash are the same address")
assertTrue(results.all { it != null && it == results.first() })
}
/** Long enough for the whole burst to pile up behind the first request. */
const val HOLD_MS = 200L
}
}
@@ -21,6 +21,8 @@
package com.vitorpamplona.quartz.nip57Zaps.validate
import com.vitorpamplona.quartz.utils.cache.ConcurrentLruCache
import kotlinx.coroutines.CompletableDeferred
import java.util.concurrent.ConcurrentHashMap
/**
* Process-wide cache of LNURL-pay endpoint metadata, keyed by the canonical
@@ -44,6 +46,12 @@ object LnurlEndpointCache {
// the previous LinkedHashMap-based behaviour where only put reordered.
private val cache = ConcurrentLruCache<String, LnurlEndpointInfo>(MAX_ENTRIES)
// Fetches in progress, keyed exactly as [cache] is. The two must agree on
// what counts as "the same address", which is why they live together: a
// flight map that keyed URLs differently would silently stop deduplicating
// the moment this object's canonicalisation changed, and nothing would fail.
private val inFlight = ConcurrentHashMap<String, CompletableDeferred<LnurlEndpointInfo?>>()
fun get(url: String): LnurlEndpointInfo? = cache.get(LnurlForm.normalizeUrl(url))
fun put(
@@ -53,8 +61,58 @@ object LnurlEndpointCache {
cache.put(LnurlForm.normalizeUrl(url), info)
}
/**
* Returns the cached entry for [url], or runs [fetch] — exactly once, however
* many callers ask at once — and caches what it returns.
*
* A zap-receipt burst asks once per receipt: twenty receipts for one
* lightning address arrive together, all miss, and without this they would
* be twenty requests to a stranger's unthrottled `/.well-known/lnurlp/`
* endpoint for one user action. The first caller fetches; the rest await it.
*
* A [fetch] returning null is not cached, so a provider having a bad minute
* is retried by the next caller rather than remembered as unresolvable.
*
* [fetch] is expected to report failure by returning null rather than by
* throwing. Cancellation is the exception: if the coroutine that owns the
* fetch is cancelled, the CancellationException propagates to that caller
* alone, and everyone awaiting it gets null — "unresolved", the same as a
* failed fetch. That asymmetry is deliberate. Awaiters are separate
* coroutines that are still alive, and cancelling them because the caller
* that happened to win the race went away would be wrong. Any other
* throwable reaches awaiters the same way, as null, while propagating to
* the owner — so a [fetch] that throws gives the two a different answer for
* one failure. Return null instead.
*/
suspend fun getOrFetch(
url: String,
fetch: suspend (String) -> LnurlEndpointInfo?,
): LnurlEndpointInfo? {
val key = LnurlForm.normalizeUrl(url)
cache.get(key)?.let { return it }
val ours = CompletableDeferred<LnurlEndpointInfo?>()
// Whoever wins putIfAbsent owns the fetch; it returns the previous entry,
// so a non-null result means someone else is already in flight.
inFlight.putIfAbsent(key, ours)?.let { return it.await() }
var info: LnurlEndpointInfo? = null
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.
inFlight.remove(key, ours)
ours.complete(info)
}
return info
}
fun clear() {
cache.clear()
inFlight.clear()
}
internal fun size(): Int = cache.size()
@@ -20,8 +20,17 @@
*/
package com.vitorpamplona.quartz.nip57Zaps.validate
import kotlinx.coroutines.CompletableDeferred
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.async
import kotlinx.coroutines.awaitAll
import kotlinx.coroutines.coroutineScope
import kotlinx.coroutines.delay
import kotlinx.coroutines.runBlocking
import java.util.concurrent.atomic.AtomicInteger
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertNotNull
import kotlin.test.assertTrue
class LnurlEndpointCacheTest {
@@ -55,4 +64,107 @@ class LnurlEndpointCacheTest {
assertEquals(null, LnurlEndpointCache.get("https://example.com/.well-known/lnurlp/user0"))
assertTrue(LnurlEndpointCache.get("https://example.com/.well-known/lnurlp/user1000") != null)
}
// --- getOrFetch: one fetch per address, however many callers ask at once ---
private val info = LnurlEndpointInfo(nostrPubkey = "d".repeat(64), allowsNostr = true)
/**
* Counts fetches and holds each one open long enough for a burst to pile up
* behind it. [failFirst] resolves the first N calls to null, as a provider
* having a bad minute does.
*/
private class CountingFetch(
private val failFirst: Int = 0,
private val result: LnurlEndpointInfo,
) {
val calls = AtomicInteger(0)
suspend fun fetch(url: String): LnurlEndpointInfo? {
val ok = calls.incrementAndGet() > failFirst
delay(HOLD_MS)
return if (ok) result else null
}
}
/**
* Releases [n] callers at once and collects what they got.
*
* The gate matters: the claim under test is "these callers were all in flight
* together", and a straggler arriving after the winner's fetch already
* returned legitimately starts a second one. Without the gate that would be a
* slow machine failing the test rather than a regression.
*/
private suspend fun burst(
n: Int,
call: suspend (Int) -> LnurlEndpointInfo?,
): List<LnurlEndpointInfo?> =
coroutineScope {
val gate = CompletableDeferred<Unit>()
val callers =
(0 until n).map { i ->
async(Dispatchers.IO) {
gate.await()
call(i)
}
}
gate.complete(Unit)
callers.awaitAll()
}
@Test
fun `concurrent callers for one url share a single fetch`() =
runBlocking {
LnurlEndpointCache.clear()
val fetcher = CountingFetch(result = info)
val results = burst(20) { LnurlEndpointCache.getOrFetch(URL, fetcher::fetch) }
assertEquals(1, fetcher.calls.get(), "20 concurrent callers must share one fetch")
assertTrue(results.all { it == info }, "every caller gets the same result")
}
@Test
fun `a failed fetch is not cached and the next caller retries`() =
runBlocking {
LnurlEndpointCache.clear()
val fetcher = CountingFetch(failFirst = 1, result = info)
val firstBurst = burst(5) { LnurlEndpointCache.getOrFetch(URL, fetcher::fetch) }
assertEquals(1, fetcher.calls.get(), "the failing burst is one fetch, not five")
assertTrue(firstBurst.all { it == null }, "a failed fetch resolves to null, not a cached verdict")
assertNotNull(LnurlEndpointCache.getOrFetch(URL, fetcher::fetch), "the retry sees the recovered provider")
assertEquals(2, fetcher.calls.get(), "the next caller retries a failed address")
LnurlEndpointCache.getOrFetch(URL, fetcher::fetch)
assertEquals(2, fetcher.calls.get(), "and the success is cached from then on")
}
@Test
fun `two spellings of one address share a single fetch`() =
runBlocking {
LnurlEndpointCache.clear()
val fetcher = CountingFetch(result = info)
// The flight map has to canonicalise the way get/put do, or a burst
// that spells the host two ways is two stampedes instead of one.
val spellings =
listOf(
"https://example.com/.well-known/lnurlp/vitor",
"https://EXAMPLE.com/.well-known/lnurlp/vitor",
"https://example.com/.well-known/lnurlp/vitor/",
)
val results = burst(spellings.size * 4) { i -> LnurlEndpointCache.getOrFetch(spellings[i % spellings.size], fetcher::fetch) }
assertEquals(1, fetcher.calls.get(), "host case and a trailing slash are the same address")
assertTrue(results.all { it == info })
}
private companion object {
const val URL = "https://example.com/.well-known/lnurlp/vitor"
/** Long enough for a whole burst to pile up behind the first fetch. */
const val HOLD_MS = 200L
}
}