mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-06 03:38:23 +00:00
refactor(dns): one resolver everywhere — promote SurgeDns to quartz, fix its blind spots, drop CachingDns
CachingDns duplicated 20% of what the app's SurgeDns already did better (stale-while-revalidate, single-flight, jittered 24-48h positive TTL, persistence, poison filtering). Consolidate on SurgeDns and fix what the verification pass found along the way: - move SurgeDns + SurgeDnsStore (+ their 41 tests) to quartz jvmAndroid so the CLI can share them; delete CachingDns - negative TTL 10s -> 10min (class default): a dead domain burns 10-30s of getaddrinfo per re-dial and both the relay pool and a crawl re-dial dead hosts continuously; the relay pool's own backoff already reaches 5min - stop the call-failure listeners from erasing negative entries they were just written from: callFailed caused by UnknownHostException must not invalidate (it deleted every negative entry milliseconds after creation, silently defeating the negative TTL entirely). The invalidation remains for its real purpose - stale positives with rotated IPs - new SurgeDns.staleAll(): soft-expire everything WITHOUT discarding - positives serve stale and revalidate in the background on next use, negatives re-try on first touch. Wired to network-identity changes in AppModules (same trigger as Tor's onNetworkChange), replacing nothing: the full-clear invalidate() was never actually called on network change - amy: SurgeDns wired into the CLI OkHttp client, snapshot persisted at ~/.amy/shared/dns-cache.bin across runs (best-effort load/save) Verified: all SurgeDns/SurgeDnsStore tests pass in the new location incl. two new staleAll tests; :amethyst compiles (main + unit tests); cli suite green. Caveat recorded in the plan doc: behind an HTTP CONNECT proxy OkHttp never consults the client resolver (hostname goes to the proxy), so this layer pays off on direct-connect deployments only - which also means the branch's A/B numbers never depended on it. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_013zEYRGKF943RgLaHTViJaB
This commit is contained in:
@@ -66,8 +66,6 @@ import com.vitorpamplona.amethyst.service.okhttp.DualHttpClientManagerForRelays
|
||||
import com.vitorpamplona.amethyst.service.okhttp.EncryptionKeyCache
|
||||
import com.vitorpamplona.amethyst.service.okhttp.OkHttpWebSocket
|
||||
import com.vitorpamplona.amethyst.service.okhttp.OnionLocationCache
|
||||
import com.vitorpamplona.amethyst.service.okhttp.SurgeDns
|
||||
import com.vitorpamplona.amethyst.service.okhttp.SurgeDnsStore
|
||||
import com.vitorpamplona.amethyst.service.playback.diskCache.VideoCache
|
||||
import com.vitorpamplona.amethyst.service.playback.diskCache.VideoCacheFactory
|
||||
import com.vitorpamplona.amethyst.service.playback.pip.BackgroundMedia
|
||||
@@ -102,6 +100,8 @@ import com.vitorpamplona.quartz.nip01Core.relay.client.accessories.RelayOfflineT
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.client.reqs.stats.RelayReqStats
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.client.stats.RelayStats
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.commands.toClient.CachingEventDecoder
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp.SurgeDns
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp.SurgeDnsStore
|
||||
import com.vitorpamplona.quartz.nip03Timestamp.VerificationStateCache
|
||||
import com.vitorpamplona.quartz.nip03Timestamp.okhttp.OkHttpBitcoinExplorer
|
||||
import com.vitorpamplona.quartz.nip03Timestamp.ots.OtsBlockHeightCache
|
||||
@@ -247,6 +247,23 @@ class AppModules(
|
||||
// path on first lookup. Stored in cacheDir — pure perf data, OK if the OS evicts it.
|
||||
val dnsStore = SurgeDnsStore(File(appContext.safeCacheDir(), SurgeDnsStore.FILE_NAME), surgeDns)
|
||||
|
||||
// Network identity change (same trigger as Tor's onNetworkChange above), but for DNS
|
||||
// the response is deliberately SOFT: most cached answers are still correct on the new
|
||||
// network, so staleAll() keeps serving every one of them and merely re-verifies —
|
||||
// positives revalidate in the background on next use, negatives (e.g. hosts that only
|
||||
// failed because the OLD network / captive portal couldn't resolve them) get re-tried
|
||||
// on first touch instead of waiting out their TTL. Nothing is dropped, nothing blocks.
|
||||
init {
|
||||
applicationIOScope.launch {
|
||||
connManager.status
|
||||
.map { (it as? ConnectivityStatus.Active)?.networkId }
|
||||
.filterNotNull()
|
||||
.distinctUntilChanged()
|
||||
.drop(1)
|
||||
.collect { surgeDns.staleAll() }
|
||||
}
|
||||
}
|
||||
|
||||
// Shared cache populated by OnionLocationInterceptor from any HTTP/WebSocket
|
||||
// response carrying an Onion-Location header. Consulted by OnionUrlRewriteInterceptor
|
||||
// on Tor-enabled clients to transparently redirect to .onion addresses.
|
||||
|
||||
+25
@@ -20,9 +20,11 @@
|
||||
*/
|
||||
package com.vitorpamplona.amethyst.service.okhttp
|
||||
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp.SurgeDns
|
||||
import okhttp3.Call
|
||||
import okhttp3.EventListener
|
||||
import java.io.IOException
|
||||
import java.net.UnknownHostException
|
||||
|
||||
/**
|
||||
* Drops a host's [SurgeDns] entry whenever an OkHttp call to it fails outright. We hook
|
||||
@@ -31,6 +33,12 @@ import java.io.IOException
|
||||
* latter while OkHttp recovers via the next address, and we don't want to invalidate the cache
|
||||
* just because the first IP was unreachable.
|
||||
*
|
||||
* A failure caused by the DNS LOOKUP ITSELF ([UnknownHostException]) must NOT invalidate:
|
||||
* this hook exists to drop STALE POSITIVES (rotated IPs that no longer connect), but for a
|
||||
* resolution failure the freshly-written negative entry IS the useful cache state — deleting
|
||||
* it milliseconds after creation re-pays the full 10-30s resolver timeout on every re-dial
|
||||
* and silently defeats the negative TTL entirely.
|
||||
*
|
||||
* Used by the relay client. The media path uses [MediaCallEventListener], which folds the same
|
||||
* invalidation into its `finish` method alongside its existing timing logging.
|
||||
*/
|
||||
@@ -41,6 +49,7 @@ class DnsInvalidatingEventListener(
|
||||
call: Call,
|
||||
ioe: IOException,
|
||||
) {
|
||||
if (ioe.causedByUnknownHost()) return
|
||||
dns.invalidate(call.request().url.host)
|
||||
}
|
||||
|
||||
@@ -53,3 +62,19 @@ class DnsInvalidatingEventListener(
|
||||
override fun create(call: Call): EventListener = listener
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* True when [UnknownHostException] appears anywhere in the cause chain — OkHttp usually
|
||||
* throws it directly from `callFailed`, but interceptors may wrap it. Depth-capped against
|
||||
* pathological cause cycles.
|
||||
*/
|
||||
fun Throwable.causedByUnknownHost(): Boolean {
|
||||
var t: Throwable? = this
|
||||
var depth = 0
|
||||
while (t != null && depth < 8) {
|
||||
if (t is UnknownHostException) return true
|
||||
t = t.cause
|
||||
depth++
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
+1
@@ -20,6 +20,7 @@
|
||||
*/
|
||||
package com.vitorpamplona.amethyst.service.okhttp
|
||||
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp.SurgeDns
|
||||
import kotlinx.coroutines.CoroutineScope
|
||||
import kotlinx.coroutines.flow.SharingStarted
|
||||
import kotlinx.coroutines.flow.StateFlow
|
||||
|
||||
+1
@@ -20,6 +20,7 @@
|
||||
*/
|
||||
package com.vitorpamplona.amethyst.service.okhttp
|
||||
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp.SurgeDns
|
||||
import kotlinx.coroutines.CoroutineScope
|
||||
import kotlinx.coroutines.flow.SharingStarted
|
||||
import kotlinx.coroutines.flow.StateFlow
|
||||
|
||||
+5
-1
@@ -21,6 +21,7 @@
|
||||
package com.vitorpamplona.amethyst.service.okhttp
|
||||
|
||||
import com.vitorpamplona.amethyst.isDebug
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp.SurgeDns
|
||||
import com.vitorpamplona.quartz.utils.Log
|
||||
import okhttp3.Call
|
||||
import okhttp3.ConnectionPool
|
||||
@@ -129,7 +130,10 @@ class MediaCallEventListener(
|
||||
// Drop the cached entry so the next attempt re-resolves instead of trying the same
|
||||
// dead IPs for up to 24h. Per-attempt connectFailed isn't enough — a multi-A-record
|
||||
// host can have one bad IP and OkHttp will recover by trying the next one.
|
||||
if (error != null) {
|
||||
// EXCEPT when the failure is the DNS lookup itself: the freshly-written negative
|
||||
// entry IS the useful cache state, and deleting it would defeat the negative TTL
|
||||
// (see DnsInvalidatingEventListener).
|
||||
if (error != null && !error.causedByUnknownHost()) {
|
||||
dns.invalidate(host)
|
||||
}
|
||||
|
||||
|
||||
+1
@@ -24,6 +24,7 @@ import com.vitorpamplona.amethyst.service.okhttp.OkHttpClientFactoryForRelays.Co
|
||||
import com.vitorpamplona.amethyst.service.okhttp.OkHttpClientFactoryForRelays.Companion.DEFAULT_SOCKS_PORT
|
||||
import com.vitorpamplona.amethyst.service.okhttp.OkHttpClientFactoryForRelays.Companion.DEFAULT_TIMEOUT_ON_MOBILE_SECS
|
||||
import com.vitorpamplona.amethyst.service.okhttp.OkHttpClientFactoryForRelays.Companion.DEFAULT_TIMEOUT_ON_WIFI_SECS
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp.SurgeDns
|
||||
import okhttp3.ConnectionPool
|
||||
import okhttp3.Dispatcher
|
||||
import okhttp3.OkHttpClient
|
||||
|
||||
+1
@@ -20,6 +20,7 @@
|
||||
*/
|
||||
package com.vitorpamplona.amethyst.service.okhttp
|
||||
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp.SurgeDns
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp.TcpNoDelaySocketFactory
|
||||
import com.vitorpamplona.quartz.utils.Log
|
||||
import okhttp3.Dispatcher
|
||||
|
||||
@@ -243,6 +243,13 @@ class DataDir(
|
||||
*/
|
||||
val eventsDbFile: File = File(eventsDir.parentFile ?: root, "events.db")
|
||||
|
||||
/**
|
||||
* Persisted DNS cache (SurgeDns snapshot) under `<root>/shared/`, next to the
|
||||
* event store — cross-account like the resolver itself: a relay host resolved
|
||||
* by one account's crawl serves every account's next run.
|
||||
*/
|
||||
val dnsCacheFile: File = File(eventsDir.parentFile ?: root, "dns-cache.bin")
|
||||
|
||||
/**
|
||||
* Machine-level operator keys for GrapeRank trusted-assertion publishing,
|
||||
* rooted at `~/.amy/operator/` (the account root's parent) so a single
|
||||
|
||||
@@ -54,7 +54,8 @@ import com.vitorpamplona.quartz.nip01Core.relay.filters.Filter
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.RelayUrlNormalizer
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp.BasicOkHttpWebSocket
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp.CachingDns
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp.SurgeDns
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp.SurgeDnsStore
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp.TcpNoDelaySocketFactory
|
||||
import com.vitorpamplona.quartz.nip01Core.signers.NostrSigner
|
||||
import com.vitorpamplona.quartz.nip01Core.signers.NostrSignerInternal
|
||||
@@ -133,6 +134,18 @@ class Context(
|
||||
*/
|
||||
val anonymous: Boolean = false,
|
||||
) : AutoCloseable {
|
||||
// Shared resolver — the SAME SurgeDns the Android app runs (stale-while-revalidate,
|
||||
// single-flight, jittered 24-48h positive TTL, 10-min negative TTL), persisted under
|
||||
// `<data-dir>/shared/` so repeat crawls start with the whole relay universe pre-resolved:
|
||||
// every previously-seen host serves its cached answer instantly and re-verifies in the
|
||||
// background instead of paying a blocking getaddrinfo per host.
|
||||
private val surgeDns = SurgeDns()
|
||||
private val dnsStore =
|
||||
SurgeDnsStore(dataDir.dnsCacheFile, surgeDns).also {
|
||||
// Best-effort warm start (~25KB read); a corrupt/missing blob just means a cold cache.
|
||||
runCatching { it.load() }
|
||||
}
|
||||
|
||||
private val okhttp =
|
||||
OkHttpClient
|
||||
.Builder()
|
||||
@@ -159,9 +172,11 @@ class Context(
|
||||
// dispatcher thread per lookup, the JVM's own cache lasts ~30s (useless
|
||||
// over a 30-minute crawl), a dead domain burns 10-30s of resolver
|
||||
// timeouts on EVERY re-dial, and the outbox model mints hundreds of
|
||||
// per-user URLs on one host — each a fresh lookup. The cache (10-min
|
||||
// positive + negative TTL, in-flight dedup per host) collapses all of it.
|
||||
.dns(CachingDns())
|
||||
// per-user URLs on one host — each a fresh lookup. SurgeDns (shared with
|
||||
// the Android app) collapses all of it: single-flight per host,
|
||||
// stale-while-revalidate positives, 10-min negative TTL, persisted
|
||||
// across runs via [dnsStore].
|
||||
.dns(surgeDns)
|
||||
.dispatcher(
|
||||
Dispatcher().apply {
|
||||
maxRequests = maxParallelHandshakes
|
||||
@@ -1022,6 +1037,10 @@ class Context(
|
||||
override fun close() {
|
||||
// Nothing to persist for an anonymous run (no account dir to write into).
|
||||
if (!anonymous) dataDir.saveRunState(state)
|
||||
// Persist the DNS cache (shared dir, account-independent) so the next run's
|
||||
// connection storm starts with every known host pre-resolved. Best-effort:
|
||||
// a failed save must never break the run it is summarizing.
|
||||
runCatching { dnsStore.save() }
|
||||
(signer as? NostrSignerRemote)?.let {
|
||||
try {
|
||||
it.closeSubscription()
|
||||
|
||||
@@ -67,9 +67,13 @@ filters — only *when* sockets open changes):
|
||||
waits once.
|
||||
2. **Transport ceilings raised (CLI).** `Dispatcher.maxRequests` 256 → 1024
|
||||
(per-host stays 16 — politeness to path-multiplexed hosts is a *server*
|
||||
property), and a `CachingDns` (10-min positive + negative TTL, in-flight
|
||||
dedup per host) replaces `Dns.SYSTEM`. FD budget is detected at startup and
|
||||
sizes both the dispatcher and the default `preconnectCap`.
|
||||
property), and the app's `SurgeDns` (moved to quartz; stale-while-revalidate,
|
||||
single-flight, 10-min negative TTL, persisted under `~/.amy/shared/`)
|
||||
replaces `Dns.SYSTEM`. FD budget is detected at startup and sizes both the
|
||||
dispatcher and the default `preconnectCap`. Caveat: in a proxied environment
|
||||
(HTTP CONNECT), OkHttp sends the hostname to the proxy and never consults
|
||||
the client-side resolver — so the DNS layer only pays off on direct-connect
|
||||
deployments, and none of the A/B numbers below depend on it.
|
||||
3. **`amy graperank probe` — the relay census.** Reads the full known relay
|
||||
universe (every relay in stored kind:10002s + everything in the
|
||||
reachability cache), mass-connects it in waves with a never-matching REQ,
|
||||
|
||||
-107
@@ -1,107 +0,0 @@
|
||||
/*
|
||||
* Copyright (c) 2025 Vitor Pamplona
|
||||
*
|
||||
* Permission is hereby granted, free of charge, to any person obtaining a copy of
|
||||
* this software and associated documentation files (the "Software"), to deal in
|
||||
* the Software without restriction, including without limitation the rights to use,
|
||||
* copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the
|
||||
* Software, and to permit persons to whom the Software is furnished to do so,
|
||||
* subject to the following conditions:
|
||||
*
|
||||
* The above copyright notice and this permission notice shall be included in all
|
||||
* copies or substantial portions of the Software.
|
||||
*
|
||||
* THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
|
||||
* IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS
|
||||
* FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR
|
||||
* COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN
|
||||
* AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION
|
||||
* WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE.
|
||||
*/
|
||||
package com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp
|
||||
|
||||
import okhttp3.Dns
|
||||
import java.net.InetAddress
|
||||
import java.net.UnknownHostException
|
||||
import java.util.concurrent.CompletableFuture
|
||||
import java.util.concurrent.ConcurrentHashMap
|
||||
|
||||
/**
|
||||
* An in-process DNS cache for OkHttp that makes mass parallel connections viable:
|
||||
*
|
||||
* - **Positive cache** ([positiveTtlMs]): the JVM's own `InetAddress` cache holds
|
||||
* entries for ~30s — useless across a 30-minute crawl that re-dials the same
|
||||
* hosts every round. Resolved addresses are kept for the TTL so reconnects and
|
||||
* the outbox model's many per-user path URLs on one host
|
||||
* (`wss://filter.example/npubA`, `/npubB`, …) resolve instantly.
|
||||
* - **Negative cache** ([negativeTtlMs]): a dead domain is the single most
|
||||
* expensive lookup — glibc `getaddrinfo` retries nameservers for 10–30s while
|
||||
* holding an OkHttp dispatcher thread — and a crawl re-dials dead relays from
|
||||
* every straggler's outbox. The first [UnknownHostException] is remembered and
|
||||
* re-thrown immediately for the TTL, so later dials fail in microseconds.
|
||||
* - **In-flight dedup**: concurrent lookups of the same host (a connect storm
|
||||
* dialing hundreds of URLs on one authority at once) collapse onto a single
|
||||
* delegate call; the rest wait on its [CompletableFuture] instead of stacking
|
||||
* N identical blocking `getaddrinfo` calls.
|
||||
*
|
||||
* Only [UnknownHostException] is negative-cached — a resolver that *errors*
|
||||
* (interrupted, SecurityException, …) is not proof the name is bad, so those
|
||||
* propagate uncached. Entries are evicted lazily on the next lookup after
|
||||
* expiry; the map is bounded by the distinct-host universe (a few thousand for
|
||||
* a full crawl), so no active eviction is needed.
|
||||
*/
|
||||
class CachingDns(
|
||||
private val delegate: Dns = Dns.SYSTEM,
|
||||
private val positiveTtlMs: Long = 10 * 60_000L,
|
||||
private val negativeTtlMs: Long = 10 * 60_000L,
|
||||
private val nowMs: () -> Long = System::currentTimeMillis,
|
||||
) : Dns {
|
||||
private class Entry(
|
||||
/** Resolved addresses, or null for a cached resolution failure. */
|
||||
val addresses: List<InetAddress>?,
|
||||
val expiresAtMs: Long,
|
||||
)
|
||||
|
||||
private val cache = ConcurrentHashMap<String, Entry>()
|
||||
private val inFlight = ConcurrentHashMap<String, CompletableFuture<List<InetAddress>>>()
|
||||
|
||||
override fun lookup(hostname: String): List<InetAddress> {
|
||||
val hit = cache[hostname]
|
||||
if (hit != null) {
|
||||
if (nowMs() < hit.expiresAtMs) {
|
||||
return hit.addresses
|
||||
?: throw UnknownHostException("$hostname (cached DNS failure)")
|
||||
}
|
||||
cache.remove(hostname, hit)
|
||||
}
|
||||
|
||||
// One resolver call per host at a time: the creator runs the delegate,
|
||||
// everyone else who raced in blocks on the same future.
|
||||
val future = CompletableFuture<List<InetAddress>>()
|
||||
val existing = inFlight.putIfAbsent(hostname, future)
|
||||
if (existing != null) {
|
||||
return try {
|
||||
existing.join()
|
||||
} catch (e: Exception) {
|
||||
throw (e.cause as? UnknownHostException) ?: UnknownHostException("$hostname (concurrent lookup failed)")
|
||||
}
|
||||
}
|
||||
|
||||
try {
|
||||
val addresses = delegate.lookup(hostname)
|
||||
cache[hostname] = Entry(addresses, nowMs() + positiveTtlMs)
|
||||
future.complete(addresses)
|
||||
return addresses
|
||||
} catch (e: UnknownHostException) {
|
||||
cache[hostname] = Entry(null, nowMs() + negativeTtlMs)
|
||||
future.completeExceptionally(e)
|
||||
throw e
|
||||
} catch (e: Exception) {
|
||||
// Not proof the name is bad (interrupt, security, …) — don't cache.
|
||||
future.completeExceptionally(e)
|
||||
throw e
|
||||
} finally {
|
||||
inFlight.remove(hostname, future)
|
||||
}
|
||||
}
|
||||
}
|
||||
+38
-3
@@ -18,7 +18,7 @@
|
||||
* AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION
|
||||
* WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE.
|
||||
*/
|
||||
package com.vitorpamplona.amethyst.service.okhttp
|
||||
package com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp
|
||||
|
||||
import okhttp3.Dns
|
||||
import java.net.InetAddress
|
||||
@@ -74,7 +74,16 @@ class SurgeDns(
|
||||
private val maxEntries: Int = 2000,
|
||||
private val positiveTtlMs: Long = TimeUnit.HOURS.toMillis(24),
|
||||
private val positiveTtlJitterMs: Long = TimeUnit.HOURS.toMillis(24),
|
||||
private val negativeTtlMs: Long = TimeUnit.SECONDS.toMillis(10),
|
||||
/**
|
||||
* How long a resolution FAILURE is served from cache. 10 minutes: a dead domain's
|
||||
* `getaddrinfo` burns 10-30s of resolver timeouts, and both the app's relay pool and a
|
||||
* GrapeRank crawl re-dial dead hosts continuously — the relay pool's own reconnect
|
||||
* backoff already reaches 5 minutes, so a 10-minute DNS negative doesn't change relay
|
||||
* behavior, it only makes the retries that do happen cost microseconds. Recovery paths
|
||||
* for a WRONG negative (transient blip): [staleAll] on a network-identity change
|
||||
* re-tries it on next use, and the entry ages out naturally.
|
||||
*/
|
||||
private val negativeTtlMs: Long = TimeUnit.MINUTES.toMillis(10),
|
||||
private val refreshExecutor: Executor = DEFAULT_REFRESH_EXECUTOR,
|
||||
) : Dns {
|
||||
private val cache = ConcurrentHashMap<String, Entry>()
|
||||
@@ -205,7 +214,33 @@ class SurgeDns(
|
||||
}
|
||||
}
|
||||
|
||||
/** Drop all cached entries. Call when the network changes (e.g. WiFi <-> mobile). */
|
||||
/**
|
||||
* Soft-expire every entry WITHOUT discarding anything — the right response to a network
|
||||
* change (WiFi <-> mobile, captive portal cleared). Most cached answers are still correct
|
||||
* on the new network, so keep serving them under a "double-check" regime instead of
|
||||
* assuming either way:
|
||||
*
|
||||
* - positives fall into the stale-while-revalidate path on their next use: the old
|
||||
* answer is served instantly and a background refresh re-verifies it (and demotes it
|
||||
* to negative if the new network really can't resolve it);
|
||||
* - negatives are re-tried synchronously on their next use, so a host that only failed
|
||||
* because of the OLD network recovers on first touch instead of waiting out its TTL.
|
||||
*
|
||||
* Contrast [invalidate], which throws every answer away and re-pays a blocking
|
||||
* `getaddrinfo` per host — a self-inflicted cold start on what is usually a healthy cache.
|
||||
*/
|
||||
fun staleAll() {
|
||||
val now = System.currentTimeMillis()
|
||||
for ((host, entry) in cache) {
|
||||
if (entry.expiresAtMillis > now) {
|
||||
// replace(k, old, new): a concurrent fresh resolution wins over the staling.
|
||||
cache.replace(host, entry, Entry(entry.addresses, now))
|
||||
}
|
||||
}
|
||||
dirty.set(true)
|
||||
}
|
||||
|
||||
/** Drop all cached entries. Prefer [staleAll] for network changes — this forces a blocking re-resolve per host. */
|
||||
fun invalidate() {
|
||||
cache.clear()
|
||||
dirty.set(true)
|
||||
+1
-1
@@ -18,7 +18,7 @@
|
||||
* AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION
|
||||
* WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE.
|
||||
*/
|
||||
package com.vitorpamplona.amethyst.service.okhttp
|
||||
package com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp
|
||||
|
||||
import com.vitorpamplona.quartz.utils.Log
|
||||
import java.io.BufferedInputStream
|
||||
-125
@@ -1,125 +0,0 @@
|
||||
/*
|
||||
* Copyright (c) 2025 Vitor Pamplona
|
||||
*
|
||||
* Permission is hereby granted, free of charge, to any person obtaining a copy of
|
||||
* this software and associated documentation files (the "Software"), to deal in
|
||||
* the Software without restriction, including without limitation the rights to use,
|
||||
* copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the
|
||||
* Software, and to permit persons to whom the Software is furnished to do so,
|
||||
* subject to the following conditions:
|
||||
*
|
||||
* The above copyright notice and this permission notice shall be included in all
|
||||
* copies or substantial portions of the Software.
|
||||
*
|
||||
* THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
|
||||
* IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS
|
||||
* FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR
|
||||
* COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN
|
||||
* AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION
|
||||
* WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE.
|
||||
*/
|
||||
package com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp
|
||||
|
||||
import okhttp3.Dns
|
||||
import java.net.InetAddress
|
||||
import java.net.UnknownHostException
|
||||
import java.util.concurrent.CountDownLatch
|
||||
import java.util.concurrent.Executors
|
||||
import java.util.concurrent.TimeUnit
|
||||
import java.util.concurrent.atomic.AtomicInteger
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertEquals
|
||||
import kotlin.test.assertFailsWith
|
||||
import kotlin.test.assertTrue
|
||||
|
||||
class CachingDnsTest {
|
||||
private val addr = listOf(InetAddress.getByAddress("good.example", byteArrayOf(10, 0, 0, 1)))
|
||||
|
||||
private class CountingDns(
|
||||
val onLookup: (String) -> List<InetAddress>,
|
||||
) : Dns {
|
||||
val calls = AtomicInteger(0)
|
||||
|
||||
override fun lookup(hostname: String): List<InetAddress> {
|
||||
calls.incrementAndGet()
|
||||
return onLookup(hostname)
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
fun cachesPositiveLookupsUntilTtl() {
|
||||
var now = 0L
|
||||
val delegate = CountingDns { addr }
|
||||
val dns = CachingDns(delegate, positiveTtlMs = 1000, negativeTtlMs = 1000, nowMs = { now })
|
||||
|
||||
assertEquals(addr, dns.lookup("good.example"))
|
||||
assertEquals(addr, dns.lookup("good.example"))
|
||||
assertEquals(1, delegate.calls.get(), "second lookup within TTL must be served from cache")
|
||||
|
||||
now = 1001
|
||||
assertEquals(addr, dns.lookup("good.example"))
|
||||
assertEquals(2, delegate.calls.get(), "expired entry must re-resolve")
|
||||
}
|
||||
|
||||
@Test
|
||||
fun cachesUnknownHostFailuresUntilTtl() {
|
||||
var now = 0L
|
||||
val delegate = CountingDns { throw UnknownHostException(it) }
|
||||
val dns = CachingDns(delegate, positiveTtlMs = 1000, negativeTtlMs = 1000, nowMs = { now })
|
||||
|
||||
assertFailsWith<UnknownHostException> { dns.lookup("dead.example") }
|
||||
assertFailsWith<UnknownHostException> { dns.lookup("dead.example") }
|
||||
assertEquals(1, delegate.calls.get(), "second failing lookup within TTL must be the cached failure")
|
||||
|
||||
now = 1001
|
||||
assertFailsWith<UnknownHostException> { dns.lookup("dead.example") }
|
||||
assertEquals(2, delegate.calls.get(), "expired negative entry must re-resolve")
|
||||
}
|
||||
|
||||
@Test
|
||||
fun nonUnknownHostErrorsAreNotCached() {
|
||||
val delegate = CountingDns { throw RuntimeException("resolver interrupted") }
|
||||
val dns = CachingDns(delegate, nowMs = { 0 })
|
||||
|
||||
assertFailsWith<RuntimeException> { dns.lookup("flaky.example") }
|
||||
assertFailsWith<RuntimeException> { dns.lookup("flaky.example") }
|
||||
assertEquals(2, delegate.calls.get(), "a transient resolver error must not be negative-cached")
|
||||
}
|
||||
|
||||
@Test
|
||||
fun concurrentLookupsOfSameHostCollapseToOneDelegateCall() {
|
||||
val started = CountDownLatch(1)
|
||||
val release = CountDownLatch(1)
|
||||
val delegate =
|
||||
CountingDns {
|
||||
started.countDown()
|
||||
release.await(5, TimeUnit.SECONDS)
|
||||
addr
|
||||
}
|
||||
val dns = CachingDns(delegate, nowMs = { 0 })
|
||||
val pool = Executors.newFixedThreadPool(8)
|
||||
try {
|
||||
val results =
|
||||
(1..8).map {
|
||||
pool.submit<List<InetAddress>> {
|
||||
if (it == 1) {
|
||||
dns.lookup("busy.example")
|
||||
} else {
|
||||
// Wait until the first lookup is inside the delegate so the
|
||||
// rest genuinely race against an in-flight resolution.
|
||||
started.await(5, TimeUnit.SECONDS)
|
||||
dns.lookup("busy.example")
|
||||
}
|
||||
}
|
||||
}
|
||||
// Give the racers a moment to pile onto the in-flight future, then release.
|
||||
started.await(5, TimeUnit.SECONDS)
|
||||
Thread.sleep(50)
|
||||
release.countDown()
|
||||
for (f in results) assertEquals(addr, f.get(5, TimeUnit.SECONDS))
|
||||
assertTrue(delegate.calls.get() <= 2, "concurrent lookups should collapse (got ${delegate.calls.get()} delegate calls)")
|
||||
} finally {
|
||||
pool.shutdownNow()
|
||||
}
|
||||
}
|
||||
}
|
||||
+1
-1
@@ -18,7 +18,7 @@
|
||||
* AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION
|
||||
* WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE.
|
||||
*/
|
||||
package com.vitorpamplona.amethyst.service.okhttp
|
||||
package com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp
|
||||
|
||||
import okhttp3.Dns
|
||||
import org.junit.Assert.assertArrayEquals
|
||||
+40
-1
@@ -18,7 +18,7 @@
|
||||
* AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION
|
||||
* WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE.
|
||||
*/
|
||||
package com.vitorpamplona.amethyst.service.okhttp
|
||||
package com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp
|
||||
|
||||
import okhttp3.Dns
|
||||
import org.junit.Assert.assertArrayEquals
|
||||
@@ -93,6 +93,45 @@ class SurgeDnsTest {
|
||||
assertEquals(1, upstream.calls("missing.example"))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `staleAll serves stale positives and revalidates on next use`() {
|
||||
val upstream = CountingDns(mapOf("a.example" to listOf(ip("1.2.3.4"))))
|
||||
val syncRefresh = Executor { it.run() }
|
||||
val dns = SurgeDns(delegate = upstream, refreshExecutor = syncRefresh)
|
||||
|
||||
dns.lookup("a.example")
|
||||
assertEquals(1, upstream.calls("a.example"))
|
||||
|
||||
dns.staleAll()
|
||||
|
||||
// Nothing was dropped: the answer still comes straight from cache — but the
|
||||
// touch double-checks upstream (background refresh, synchronous here).
|
||||
assertEquals(listOf(ip("1.2.3.4")), dns.lookup("a.example"))
|
||||
assertEquals(2, upstream.calls("a.example"))
|
||||
|
||||
// The refresh re-armed the TTL: back to a plain cache hit.
|
||||
dns.lookup("a.example")
|
||||
assertEquals(2, upstream.calls("a.example"))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `staleAll retries negatives on next use instead of waiting out the TTL`() {
|
||||
val upstream = CountingDns(emptyMap())
|
||||
val dns = SurgeDns(delegate = upstream)
|
||||
|
||||
assertThrows(UnknownHostException::class.java) { dns.lookup("missing.example") }
|
||||
assertThrows(UnknownHostException::class.java) { dns.lookup("missing.example") }
|
||||
assertEquals(1, upstream.calls("missing.example"))
|
||||
|
||||
dns.staleAll()
|
||||
|
||||
// The negative no longer short-circuits: the next use goes back upstream (sync
|
||||
// path — negatives never serve stale), so a host that only failed because of
|
||||
// the OLD network recovers on first touch after a network change.
|
||||
assertThrows(UnknownHostException::class.java) { dns.lookup("missing.example") }
|
||||
assertEquals(2, upstream.calls("missing.example"))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `expired positive entry serves stale and refreshes in background`() {
|
||||
val upstream = CountingDns(mapOf("a.example" to listOf(ip("1.2.3.4"))))
|
||||
Reference in New Issue
Block a user