mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-05 19:28:25 +00:00
refactor: rename accessory timeoutMs to idleTimeoutMs
Every wait in the accessories package is an idle window measured from the
relay's most recent progress, so the parameter now says so. The name is
the contract: a caller reading timeoutMs reasonably expects a deadline,
which is exactly the misreading that made fetchAllPages' hard per-page
cap look correct for so long.
Renamed across the public surface -- fetchAll (7 overloads),
fetchAllWithHooks, fetchAllPages (2), fetchAllPagesFromPool,
fetchAllPagesFromPoolWithHooks, fetchFirst, count (2), countMerged -- plus
the quartz wrappers whose own parameter is a pure pass-through of that
window (KeyPackageFetcher, RecipientRelayFetcher, FollowerCrawler.Config)
and every call site across commons, cli, desktopApp and amethyst.
Deliberately NOT renamed, because these are genuine wall-clock bounds and
the differing name is the tell:
- publishAndConfirm's timeoutInSeconds -- one fixed window to collect
the OKs, a bounded confirmation round-trip rather than a stream.
- GrapeRankCrawler.Config.timeoutMs -- a hard per-drain gate
(withTimeoutOrNull(config.timeoutMs)) that also drives parking.
- Context.awaitReply's timeoutMs, Context.syncIncoming, and the
non-accessory app-layer helpers (RelayProber, FeedMetadataCoordinator,
RelayAuthPromptBus).
This is a source-breaking change for named-argument callers of quartz.
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01KZe1Pq2ejoHehdufhPDRgb
This commit is contained in:
@@ -971,7 +971,7 @@ class AccountConcordActions(
|
||||
val filter = Filter(kinds = listOf(ConcordCommunityListEvent.KIND), authors = listOf(account.signer.pubKey))
|
||||
// Stock relays like relay.ditto.pub can be slow (~10–20s to first response), so give
|
||||
// the fetch a generous window to drain every relay before we pick the newest copy.
|
||||
val events = account.client.fetchAll(filters = relays.associateWith { listOf(filter) }, timeoutMs = 30_000L)
|
||||
val events = account.client.fetchAll(filters = relays.associateWith { listOf(filter) }, idleTimeoutMs = 30_000L)
|
||||
val newest = events.filterIsInstance<ConcordCommunityListEvent>().maxByOrNull { it.createdAt }
|
||||
val entryCount = newest?.let { runCatching { it.decrypt(account.signer).size }.getOrElse { -1 } } ?: 0
|
||||
Log.d(
|
||||
@@ -1015,7 +1015,7 @@ class AccountConcordActions(
|
||||
}
|
||||
if (filters.isEmpty()) return
|
||||
val byRelay = filters.groupBy { it.relay }.mapValues { (_, group) -> group.map { it.filter } }
|
||||
account.client.fetchAll(filters = byRelay, timeoutMs = 20_000L)
|
||||
account.client.fetchAll(filters = byRelay, idleTimeoutMs = 20_000L)
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -133,7 +133,7 @@ class AccountRelayGroupActions(
|
||||
if (channelId == null && results.values.any { !it.accepted && it.message.contains("auth-required", ignoreCase = true) }) {
|
||||
account.client.fetchAllWithHooks(
|
||||
filters = mapOf(relay to listOf(Filter(kinds = listOf(DmOpenEvent.KIND), limit = 1))),
|
||||
timeoutMs = 8_000,
|
||||
idleTimeoutMs = 8_000,
|
||||
pendingOnAuthRequired = true,
|
||||
) { _, _ -> false }
|
||||
results = account.client.publishAndCollectResults(signed, setOf(relay))
|
||||
|
||||
+1
-1
@@ -255,7 +255,7 @@ class AccountNappletGateways(
|
||||
emptyList()
|
||||
} else {
|
||||
runCatching {
|
||||
account.client.fetchAll(filters = relays.associateWith { filters }, timeoutMs = QUERY_TIMEOUT.inWholeMilliseconds)
|
||||
account.client.fetchAll(filters = relays.associateWith { filters }, idleTimeoutMs = QUERY_TIMEOUT.inWholeMilliseconds)
|
||||
}.getOrDefault(emptyList())
|
||||
}
|
||||
val fromCache = filters.flatMap { filter -> account.cache.filter(filter).mapNotNull { it.event } }
|
||||
|
||||
+1
-1
@@ -144,7 +144,7 @@ class NappletResourceFetcher(
|
||||
val relays = account.homeRelays.flow.value
|
||||
if (relays.isEmpty()) return null
|
||||
return runCatching {
|
||||
account.client.fetchAll(filters = relays.associateWith { listOf(filter) }, timeoutMs = NOSTR_FETCH_TIMEOUT_MS)
|
||||
account.client.fetchAll(filters = relays.associateWith { listOf(filter) }, idleTimeoutMs = NOSTR_FETCH_TIMEOUT_MS)
|
||||
}.getOrDefault(emptyList())
|
||||
.maxByOrNull { it.createdAt }
|
||||
}
|
||||
|
||||
+1
-1
@@ -153,7 +153,7 @@ class AgentConsoleViewModel : ViewModel() {
|
||||
// (pendingOnAuthRequired) so it authenticates on the `auth-required` CLOSED and retries.
|
||||
account.client.fetchAllWithHooks(
|
||||
filters = relays.associateWith { filters },
|
||||
timeoutMs = 8_000,
|
||||
idleTimeoutMs = 8_000,
|
||||
pendingOnAuthRequired = true,
|
||||
) { _, _ -> false }
|
||||
}
|
||||
|
||||
+2
-2
@@ -97,7 +97,7 @@ private suspend fun runBuzzDmDiscovery(
|
||||
// rather than returning empty.
|
||||
account.client.fetchAllWithHooks(
|
||||
filters = relays.associateWith { discoveryFilters },
|
||||
timeoutMs = 8_000,
|
||||
idleTimeoutMs = 8_000,
|
||||
pendingOnAuthRequired = true,
|
||||
) { relay, event ->
|
||||
(event as? MemberAddedNotificationEvent)?.let { recordDiscovery(me, it, relay) }
|
||||
@@ -135,7 +135,7 @@ private suspend fun fetchDmMetadata(
|
||||
.groupBy({ it.value }, { it.key })
|
||||
.mapValues { (_, ids) -> listOf(Filter(kinds = RELAY_GROUP_METADATA_KINDS, tags = mapOf("d" to ids))) }
|
||||
if (byRelay.isEmpty()) return
|
||||
account.client.fetchAllWithHooks(filters = byRelay, timeoutMs = 8_000, pendingOnAuthRequired = true) { _, _ -> false }
|
||||
account.client.fetchAllWithHooks(filters = byRelay, idleTimeoutMs = 8_000, pendingOnAuthRequired = true) { _, _ -> false }
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
+2
-2
@@ -200,7 +200,7 @@ class BuzzDmListViewModel : ViewModel() {
|
||||
)
|
||||
account.client.fetchAllWithHooks(
|
||||
filters = relays.associateWith { filters },
|
||||
timeoutMs = 8_000,
|
||||
idleTimeoutMs = 8_000,
|
||||
pendingOnAuthRequired = true,
|
||||
) { relay, event ->
|
||||
(event as? MemberAddedNotificationEvent)?.channel()?.let { memberChannels[it] = relay }
|
||||
@@ -215,7 +215,7 @@ class BuzzDmListViewModel : ViewModel() {
|
||||
.groupBy({ it.value }, { it.key })
|
||||
.mapValues { (_, ids) -> listOf(Filter(kinds = RELAY_GROUP_METADATA_KINDS, tags = mapOf("d" to ids))) }
|
||||
if (byRelay.isEmpty()) return
|
||||
account.client.fetchAllWithHooks(filters = byRelay, timeoutMs = 8_000, pendingOnAuthRequired = true) { _, _ -> false }
|
||||
account.client.fetchAllWithHooks(filters = byRelay, idleTimeoutMs = 8_000, pendingOnAuthRequired = true) { _, _ -> false }
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
+2
-2
@@ -149,7 +149,7 @@ class BuzzRelayImportViewModel : ViewModel() {
|
||||
),
|
||||
),
|
||||
),
|
||||
timeoutMs = 8_000,
|
||||
idleTimeoutMs = 8_000,
|
||||
pendingOnAuthRequired = true,
|
||||
) { _, event ->
|
||||
(event as? MemberAddedNotificationEvent)?.channel()?.let { channelIds.add(it) }
|
||||
@@ -172,7 +172,7 @@ class BuzzRelayImportViewModel : ViewModel() {
|
||||
Filter(kinds = listOf(SystemMessageEvent.KIND), tags = mapOf("h" to channelIds.toList())),
|
||||
),
|
||||
),
|
||||
timeoutMs = 8_000,
|
||||
idleTimeoutMs = 8_000,
|
||||
pendingOnAuthRequired = true,
|
||||
) { _, _ -> false }
|
||||
}
|
||||
|
||||
+1
-1
@@ -502,7 +502,7 @@ class EventSync(
|
||||
try {
|
||||
client.fetchAllPagesFromPool(
|
||||
filters = perRelayFilters,
|
||||
timeoutMs = RELAY_TIMEOUT_MS,
|
||||
idleTimeoutMs = RELAY_TIMEOUT_MS,
|
||||
maxConcurrentRelays = MAX_CONCURRENT_RELAYS,
|
||||
onNewPage = { until, sourceRelay ->
|
||||
_liveActivity.value.runningRelays[sourceRelay]
|
||||
|
||||
+1
-1
@@ -199,7 +199,7 @@ class CashuWalletDiscovery(
|
||||
fetchAllPages(
|
||||
relay = relay,
|
||||
filters = filters,
|
||||
timeoutMs = RELAY_TIMEOUT_MS,
|
||||
idleTimeoutMs = RELAY_TIMEOUT_MS,
|
||||
onEvent = onEvent,
|
||||
)
|
||||
}.onFailure {
|
||||
|
||||
+15
-15
@@ -153,12 +153,12 @@ class AmethystAppFunctions {
|
||||
|
||||
// Quartz's INostrClient.fetchAll handles subscribe → drain on
|
||||
// EOSE/closed/cannot-connect → unsubscribe → dedup by id → sort
|
||||
// newest-first. Wraps everything in a withTimeoutOrNull(timeoutMs)
|
||||
// newest-first. Wraps everything in a withTimeoutOrNull(idleTimeoutMs)
|
||||
// so a slow relay can't stall the dispatch.
|
||||
val events =
|
||||
client.fetchAll(
|
||||
filters = relays.associateWith { listOf(filter) },
|
||||
timeoutMs = GEMINI_FETCH_TIMEOUT_MS,
|
||||
idleTimeoutMs = GEMINI_FETCH_TIMEOUT_MS,
|
||||
)
|
||||
|
||||
val candidates =
|
||||
@@ -397,7 +397,7 @@ class AmethystAppFunctions {
|
||||
return Amethyst.instance.client
|
||||
.fetchAll(
|
||||
filters = relays.associateWith { listOf(filter) },
|
||||
timeoutMs = GEMINI_FETCH_TIMEOUT_MS,
|
||||
idleTimeoutMs = GEMINI_FETCH_TIMEOUT_MS,
|
||||
).mapNotNull { it as? TextNoteEvent }
|
||||
.take(limit)
|
||||
}
|
||||
@@ -449,7 +449,7 @@ class AmethystAppFunctions {
|
||||
val events =
|
||||
client.fetchAll(
|
||||
filters = relays.associateWith { listOf(filter) },
|
||||
timeoutMs = GEMINI_FETCH_TIMEOUT_MS,
|
||||
idleTimeoutMs = GEMINI_FETCH_TIMEOUT_MS,
|
||||
)
|
||||
|
||||
val hits =
|
||||
@@ -521,7 +521,7 @@ class AmethystAppFunctions {
|
||||
client
|
||||
.fetchAll(
|
||||
filters = relays.associateWith { listOf(filter) },
|
||||
timeoutMs = GEMINI_FETCH_TIMEOUT_MS,
|
||||
idleTimeoutMs = GEMINI_FETCH_TIMEOUT_MS,
|
||||
).mapNotNull { it as? MetadataEvent }
|
||||
.filter { it.pubKey == pubkey }
|
||||
.maxByOrNull { it.createdAt }
|
||||
@@ -569,7 +569,7 @@ class AmethystAppFunctions {
|
||||
val events =
|
||||
client.fetchAll(
|
||||
filters = relays.associateWith { listOf(filter) },
|
||||
timeoutMs = GEMINI_FETCH_TIMEOUT_MS,
|
||||
idleTimeoutMs = GEMINI_FETCH_TIMEOUT_MS,
|
||||
)
|
||||
|
||||
val hits =
|
||||
@@ -643,7 +643,7 @@ class AmethystAppFunctions {
|
||||
val events =
|
||||
client.fetchAll(
|
||||
filters = relays.associateWith { listOf(filter) },
|
||||
timeoutMs = GEMINI_FETCH_TIMEOUT_MS,
|
||||
idleTimeoutMs = GEMINI_FETCH_TIMEOUT_MS,
|
||||
)
|
||||
|
||||
val hits =
|
||||
@@ -686,7 +686,7 @@ class AmethystAppFunctions {
|
||||
val events =
|
||||
client.fetchAll(
|
||||
filters = relays.associateWith { listOf(filter) },
|
||||
timeoutMs = GEMINI_FETCH_TIMEOUT_MS,
|
||||
idleTimeoutMs = GEMINI_FETCH_TIMEOUT_MS,
|
||||
)
|
||||
|
||||
val hits =
|
||||
@@ -733,7 +733,7 @@ class AmethystAppFunctions {
|
||||
val events =
|
||||
client.fetchAll(
|
||||
filters = relays.associateWith { listOf(filter) },
|
||||
timeoutMs = GEMINI_FETCH_TIMEOUT_MS,
|
||||
idleTimeoutMs = GEMINI_FETCH_TIMEOUT_MS,
|
||||
)
|
||||
|
||||
val hits =
|
||||
@@ -785,7 +785,7 @@ class AmethystAppFunctions {
|
||||
val events =
|
||||
client.fetchAll(
|
||||
filters = relays.associateWith { listOf(filter) },
|
||||
timeoutMs = GEMINI_FETCH_TIMEOUT_MS,
|
||||
idleTimeoutMs = GEMINI_FETCH_TIMEOUT_MS,
|
||||
)
|
||||
|
||||
val receipts = events.mapNotNull { it as? LnZapEvent }
|
||||
@@ -880,7 +880,7 @@ class AmethystAppFunctions {
|
||||
client
|
||||
.fetchAll(
|
||||
filters = relays.associateWith { listOf(filter) },
|
||||
timeoutMs = GEMINI_FETCH_TIMEOUT_MS,
|
||||
idleTimeoutMs = GEMINI_FETCH_TIMEOUT_MS,
|
||||
).mapNotNull { it as? GiftWrapEvent }
|
||||
|
||||
val seen = HashSet<HexKey>()
|
||||
@@ -947,7 +947,7 @@ class AmethystAppFunctions {
|
||||
val events =
|
||||
client.fetchAll(
|
||||
filters = relays.associateWith { listOf(filter) },
|
||||
timeoutMs = GEMINI_FETCH_TIMEOUT_MS,
|
||||
idleTimeoutMs = GEMINI_FETCH_TIMEOUT_MS,
|
||||
)
|
||||
|
||||
val hits =
|
||||
@@ -991,7 +991,7 @@ class AmethystAppFunctions {
|
||||
val events =
|
||||
client.fetchAll(
|
||||
filters = relays.associateWith { listOf(filter) },
|
||||
timeoutMs = GEMINI_FETCH_TIMEOUT_MS,
|
||||
idleTimeoutMs = GEMINI_FETCH_TIMEOUT_MS,
|
||||
)
|
||||
|
||||
val streams =
|
||||
@@ -1691,7 +1691,7 @@ class AmethystAppFunctions {
|
||||
return client
|
||||
.fetchAll(
|
||||
filters = relays.associateWith { listOf(filter) },
|
||||
timeoutMs = GEMINI_FETCH_TIMEOUT_MS,
|
||||
idleTimeoutMs = GEMINI_FETCH_TIMEOUT_MS,
|
||||
).mapNotNull { it as? MetadataEvent }
|
||||
.maxByOrNull { it.createdAt }
|
||||
?.contactMetaData()
|
||||
@@ -2027,7 +2027,7 @@ class AmethystAppFunctions {
|
||||
val events =
|
||||
client.fetchAll(
|
||||
filters = relays.associateWith { listOf(filter) },
|
||||
timeoutMs = GEMINI_FETCH_TIMEOUT_MS,
|
||||
idleTimeoutMs = GEMINI_FETCH_TIMEOUT_MS,
|
||||
)
|
||||
|
||||
val hits =
|
||||
|
||||
@@ -547,7 +547,7 @@ class Context(
|
||||
if (needAuth.isEmpty()) break
|
||||
// A cheap REQ whose only purpose is to force the AUTH handshake to completion.
|
||||
val warmFilter = listOf(Filter(kinds = listOf(event.kind), limit = 1))
|
||||
drain(needAuth.associateWith { warmFilter }, timeoutMs = 8_000, pendingOnAuthRequired = true)
|
||||
drain(needAuth.associateWith { warmFilter }, idleTimeoutMs = 8_000, pendingOnAuthRequired = true)
|
||||
results = results + client.publishAndCollectResults(event, needAuth, timeoutSecs)
|
||||
attempt++
|
||||
}
|
||||
@@ -565,7 +565,7 @@ class Context(
|
||||
* When [deadOut] is provided, every relay that reported it could not be
|
||||
* connected to (`onCannotConnect`) is added to it, so callers can prune
|
||||
* proven-dead relays from future routing instead of paying the full
|
||||
* [timeoutMs] on them again. Slow-but-connected relays are NOT reported —
|
||||
* [idleTimeoutMs] on them again. Slow-but-connected relays are NOT reported —
|
||||
* only hard connect failures, so a temporarily-busy relay isn't discarded.
|
||||
*
|
||||
* With [pendingOnAuthRequired], a relay that refuses the REQ with an
|
||||
@@ -573,24 +573,24 @@ class Context(
|
||||
* NIP-42 responder answers the challenge and the client re-fires this same
|
||||
* subscription (`syncFilters`), so the post-auth events are collected instead of
|
||||
* returning empty. If auth never satisfies it, the relay simply falls through to
|
||||
* the [timeoutMs]. Needed for Concord planes, whose kind-1059 wraps are served
|
||||
* the [idleTimeoutMs]. Needed for Concord planes, whose kind-1059 wraps are served
|
||||
* only to a connection authenticated as the derived stream key.
|
||||
*/
|
||||
suspend fun drain(
|
||||
filters: Map<NormalizedRelayUrl, List<Filter>>,
|
||||
timeoutMs: Long = 8_000,
|
||||
idleTimeoutMs: Long = 8_000,
|
||||
diagnoseSlow: Boolean = false,
|
||||
deadOut: MutableMap<NormalizedRelayUrl, DrainFailure>? = null,
|
||||
pendingOnAuthRequired: Boolean = false,
|
||||
): List<Pair<NormalizedRelayUrl, Event>> =
|
||||
client.fetchAllWithHooks(
|
||||
filters = filters,
|
||||
timeoutMs = timeoutMs,
|
||||
idleTimeoutMs = idleTimeoutMs,
|
||||
pendingOnAuthRequired = pendingOnAuthRequired,
|
||||
deadOut = deadOut,
|
||||
onTimeout =
|
||||
if (diagnoseSlow) {
|
||||
{ stalled, doneReasons, collected -> logSlowDrain(timeoutMs, stalled, doneReasons, collected) }
|
||||
{ stalled, doneReasons, collected -> logSlowDrain(idleTimeoutMs, stalled, doneReasons, collected) }
|
||||
} else {
|
||||
null
|
||||
},
|
||||
@@ -604,7 +604,7 @@ class Context(
|
||||
* "relay is slow" and "we never connected" are easy to tell apart.
|
||||
*/
|
||||
private fun logSlowDrain(
|
||||
timeoutMs: Long,
|
||||
idleTimeoutMs: Long,
|
||||
stalled: Set<NormalizedRelayUrl>,
|
||||
doneReasons: Map<NormalizedRelayUrl, String>,
|
||||
collected: List<Pair<NormalizedRelayUrl, Event>>,
|
||||
@@ -615,7 +615,7 @@ class Context(
|
||||
val slowDetail = stalled.take(12).joinToString(", ") { "${it.url}(${eventsPer[it] ?: 0}ev)" }
|
||||
val cannotDetail = cannot.entries.take(8).joinToString(", ") { "${it.key.url}=${it.value.removePrefix("cannot:").take(40)}" }
|
||||
System.err.println(
|
||||
"[drain] timeout ${timeoutMs}ms: ${stalled.size} slow(no EOSE), ${cannot.size} cannot-connect, ${closed.size} closed" +
|
||||
"[drain] timeout ${idleTimeoutMs}ms: ${stalled.size} slow(no EOSE), ${cannot.size} cannot-connect, ${closed.size} closed" +
|
||||
(if (slowDetail.isNotEmpty()) " | slow: $slowDetail" else "") +
|
||||
(if (cannotDetail.isNotEmpty()) " | cannot: $cannotDetail" else ""),
|
||||
)
|
||||
@@ -641,12 +641,12 @@ class Context(
|
||||
*/
|
||||
suspend fun drainAllPages(
|
||||
filters: Map<NormalizedRelayUrl, List<Filter>>,
|
||||
timeoutMs: Long = 30_000,
|
||||
idleTimeoutMs: Long = 30_000,
|
||||
maxConcurrentRelays: Int = 8,
|
||||
): List<Pair<NormalizedRelayUrl, Event>> =
|
||||
client.fetchAllPagesFromPoolWithHooks(
|
||||
filters = filters,
|
||||
timeoutMs = timeoutMs,
|
||||
idleTimeoutMs = idleTimeoutMs,
|
||||
maxConcurrentRelays = maxConcurrentRelays,
|
||||
) { _, event -> verifyAndStore(event) }
|
||||
|
||||
|
||||
@@ -112,7 +112,7 @@ object AwaitCommands {
|
||||
val event =
|
||||
ctx.client.fetchFirst(
|
||||
filters = relays.associateWith { listOf(filter) },
|
||||
timeoutMs = 3_000,
|
||||
idleTimeoutMs = 3_000,
|
||||
)
|
||||
if (event is KeyPackageEvent) {
|
||||
Output.emit(
|
||||
|
||||
@@ -357,7 +357,7 @@ object DmCommands {
|
||||
.groupBy { it.relay }
|
||||
.mapValues { (_, v) -> v.map { it.filter } }
|
||||
|
||||
val raw = ctx.drain(filters, timeoutMs = timeoutSecs * 1000)
|
||||
val raw = ctx.drain(filters, idleTimeoutMs = timeoutSecs * 1000)
|
||||
|
||||
val messages = decryptDms(ctx, raw, peerHex)
|
||||
val out =
|
||||
@@ -415,7 +415,7 @@ object DmCommands {
|
||||
.groupBy { it.relay }
|
||||
.mapValues { (_, v) -> v.map { it.filter } }
|
||||
|
||||
val raw = ctx.drain(filters, timeoutMs = 3_000)
|
||||
val raw = ctx.drain(filters, idleTimeoutMs = 3_000)
|
||||
val messages = decryptDms(ctx, raw, peerHex)
|
||||
// Match against the text body for kind:14 and against the URL
|
||||
// for kind:15 — both are exposed as `searchText` so callers
|
||||
|
||||
@@ -103,7 +103,7 @@ object GitReadCommands {
|
||||
ctx
|
||||
.drainAllPages(
|
||||
relays.associateWith { listOf(Filter(kinds = listOf(itemKind), tags = mapOf("a" to listOf(repoAddress)), limit = limit)) },
|
||||
timeoutMs = READ_TIMEOUT_MS,
|
||||
idleTimeoutMs = READ_TIMEOUT_MS,
|
||||
).asSequence()
|
||||
.map { it.second }
|
||||
.filter { it.kind == itemKind }
|
||||
@@ -160,7 +160,7 @@ object GitReadCommands {
|
||||
relays.associateWith {
|
||||
listOf(Filter(kinds = STATUS_KINDS + listOf(CommentEvent.KIND, GitReplyEvent.KIND), tags = mapOf("e" to listOf(id))))
|
||||
},
|
||||
timeoutMs = READ_TIMEOUT_MS,
|
||||
idleTimeoutMs = READ_TIMEOUT_MS,
|
||||
).map { it.second }
|
||||
.distinctBy { it.id }
|
||||
|
||||
@@ -208,7 +208,7 @@ object GitReadCommands {
|
||||
ctx
|
||||
.drainAllPages(
|
||||
relays.associateWith { listOf(Filter(kinds = STATUS_KINDS, tags = mapOf("e" to chunk))) },
|
||||
timeoutMs = READ_TIMEOUT_MS,
|
||||
idleTimeoutMs = READ_TIMEOUT_MS,
|
||||
).map { it.second }
|
||||
}.filterIsInstance<GitStatusEvent>()
|
||||
.distinctBy { it.id }
|
||||
|
||||
+1
-1
@@ -97,7 +97,7 @@ object GroupAddMemberCommand {
|
||||
client = ctx.client,
|
||||
targetPubKey = pub,
|
||||
relays = kpRelays,
|
||||
timeoutMs = 10_000,
|
||||
idleTimeoutMs = 10_000,
|
||||
)
|
||||
if (kpEvent == null) {
|
||||
report.add(mapOf("pubkey" to pub, "status" to "no_key_package"))
|
||||
|
||||
@@ -99,7 +99,7 @@ object KeyPackageCommands {
|
||||
client = ctx.client,
|
||||
targetPubKey = targetHex,
|
||||
relays = relays,
|
||||
timeoutMs = 10_000,
|
||||
idleTimeoutMs = 10_000,
|
||||
)
|
||||
if (event == null) {
|
||||
return Output.error("not_found", "no KeyPackage for $targetHex on ${relays.size} relay(s)")
|
||||
|
||||
+1
-1
@@ -287,7 +287,7 @@ object GrapeRankCrawl {
|
||||
// Default null → pull EVERY follower each relay holds; --max
|
||||
// N caps the total per relay for a quick spot check.
|
||||
maxPerRelay = args.flag("max")?.toIntOrNull(),
|
||||
timeoutMs = args.timeoutMs(15),
|
||||
idleTimeoutMs = args.timeoutMs(15),
|
||||
maxConcurrentRelays = relayConcurrency,
|
||||
insertBatchSize = args.intFlag(FLAG_INSERT_BATCH, INSERT_BATCH_DEFAULT),
|
||||
),
|
||||
|
||||
+1
-1
@@ -147,7 +147,7 @@ class DmInboxRelayResolver(
|
||||
if (writeRelays.isNotEmpty()) {
|
||||
relays =
|
||||
RecipientRelayFetcher
|
||||
.fetchRelayLists(unauthenticatedClient, pubkey, writeRelays, timeoutMs = 5_000L)
|
||||
.fetchRelayLists(unauthenticatedClient, pubkey, writeRelays, idleTimeoutMs = 5_000L)
|
||||
.dmInbox
|
||||
}
|
||||
}
|
||||
|
||||
+1
-1
@@ -221,7 +221,7 @@ class DesktopRelaySubscriptionsCoordinator(
|
||||
val events =
|
||||
client.fetchAll(
|
||||
filters = indexRelays.associateWith { listOf(filter) },
|
||||
timeoutMs = 8.seconds.inWholeMilliseconds,
|
||||
idleTimeoutMs = 8.seconds.inWholeMilliseconds,
|
||||
)
|
||||
events.forEach { consumeEvent(it, null) }
|
||||
}
|
||||
|
||||
+3
-3
@@ -75,14 +75,14 @@ class FollowerCrawler(
|
||||
* stops paging once it's reached, so a non-null value cuts the crawl short at
|
||||
* that many per relay. Leave it null for completeness; set it only to bound a
|
||||
* spot check. The per-page size is the relay's own default either way.
|
||||
* @param timeoutMs per-page EOSE timeout for a relay before its next page fires.
|
||||
* @param idleTimeoutMs per-page EOSE timeout for a relay before its next page fires.
|
||||
* @param maxConcurrentRelays how many relays page at once (a global fan-out cap).
|
||||
* @param insertBatchSize verified events group-committed per [IEventStore.batchInsert].
|
||||
*/
|
||||
class Config(
|
||||
val relays: Set<NormalizedRelayUrl>,
|
||||
val maxPerRelay: Int? = null,
|
||||
val timeoutMs: Long = 15_000,
|
||||
val idleTimeoutMs: Long = 15_000,
|
||||
val maxConcurrentRelays: Int = 16,
|
||||
val insertBatchSize: Int = 500,
|
||||
)
|
||||
@@ -162,7 +162,7 @@ class FollowerCrawler(
|
||||
|
||||
client.fetchAllPagesFromPool(
|
||||
filters = perRelay,
|
||||
timeoutMs = config.timeoutMs,
|
||||
idleTimeoutMs = config.idleTimeoutMs,
|
||||
maxConcurrentRelays = config.maxConcurrentRelays,
|
||||
onRelayComplete = { relay, total ->
|
||||
if (total > 0) log("[followers] ${relay.url}: $total kind:3 pages drained")
|
||||
|
||||
+2
-2
@@ -81,7 +81,7 @@ object RecipientRelayFetcher {
|
||||
client: INostrClient,
|
||||
pubKey: HexKey,
|
||||
seedRelays: Set<NormalizedRelayUrl>,
|
||||
timeoutMs: Long = 8_000L,
|
||||
idleTimeoutMs: Long = 8_000L,
|
||||
): Lists {
|
||||
if (seedRelays.isEmpty()) return Lists(emptyList(), emptyList(), null)
|
||||
|
||||
@@ -99,7 +99,7 @@ object RecipientRelayFetcher {
|
||||
val events =
|
||||
client.fetchAll(
|
||||
filters = seedRelays.associateWith { listOf(filter) },
|
||||
timeoutMs = timeoutMs,
|
||||
idleTimeoutMs = idleTimeoutMs,
|
||||
)
|
||||
|
||||
var dm: ChatMessageRelayListEvent? = null
|
||||
|
||||
+2
-2
@@ -75,11 +75,11 @@ object KeyPackageFetcher {
|
||||
client: INostrClient,
|
||||
targetPubKey: HexKey,
|
||||
relays: Set<NormalizedRelayUrl>,
|
||||
timeoutMs: Long = 30_000,
|
||||
idleTimeoutMs: Long = 30_000,
|
||||
): KeyPackageEvent? {
|
||||
if (relays.isEmpty()) return null
|
||||
val filter = MarmotFilters.keyPackagesByAuthor(targetPubKey)
|
||||
val events = client.fetchAll(filters = relays.associateWith { listOf(filter) }, timeoutMs = timeoutMs)
|
||||
val events = client.fetchAll(filters = relays.associateWith { listOf(filter) }, idleTimeoutMs = idleTimeoutMs)
|
||||
// fetchAll returns events sorted by created_at DESC, so the first
|
||||
// KeyPackageEvent is the most recent one any relay had.
|
||||
return events.firstNotNullOfOrNull { it as? KeyPackageEvent }
|
||||
|
||||
+5
-4
@@ -30,10 +30,11 @@ import kotlin.time.TimeSource
|
||||
* fetch/sync loops. [bump] on every sign of life from the relay; [elapsedMs] reports
|
||||
* the silence since the last bump (or since construction, before the first bump).
|
||||
*
|
||||
* This is the timeout convention for every accessory in this package: a `timeoutMs`
|
||||
* (or `idleTimeoutMs`) is an **idle window measured from the relay's most recent
|
||||
* message**, not a wall-clock deadline — an actively streaming relay is never cut
|
||||
* off mid-delivery, only one that goes silent.
|
||||
* This is the timeout convention for every accessory in this package: an
|
||||
* `idleTimeoutMs` is an **idle window measured from the relay's most recent
|
||||
* progress**, not a wall-clock deadline — an actively streaming relay is never cut
|
||||
* off mid-delivery, only one that goes silent. The name is the contract: a
|
||||
* parameter here is called `idleTimeoutMs` precisely because it is not a deadline.
|
||||
*
|
||||
* [bump] is on the per-event hot path (a connection listener may bump for every
|
||||
* message the relay sends — millions during a large download), so it must not
|
||||
|
||||
+11
-11
@@ -38,19 +38,19 @@ import kotlinx.coroutines.withTimeoutOrNull
|
||||
* Sends a NIP-45 COUNT query to a single relay and suspends until
|
||||
* the result arrives or the timeout expires.
|
||||
*
|
||||
* A COUNT exchange is a single response message, so [timeoutMs] here is
|
||||
* A COUNT exchange is a single response message, so [idleTimeoutMs] here is
|
||||
* trivially the package-wide idle-window convention (time since the most
|
||||
* recent message): no message can arrive before the one that completes it.
|
||||
*
|
||||
* @param relay Target relay to query.
|
||||
* @param filter The filter to count against.
|
||||
* @param timeoutMs How long to wait for the response (default 15 s).
|
||||
* @param idleTimeoutMs How long to wait for the response (default 15 s).
|
||||
* @return The [CountResult], or `null` on timeout.
|
||||
*/
|
||||
suspend fun INostrClient.count(
|
||||
relay: NormalizedRelayUrl,
|
||||
filter: Filter,
|
||||
timeoutMs: Long = 15_000,
|
||||
idleTimeoutMs: Long = 15_000,
|
||||
): CountResult? {
|
||||
val subId = newSubId()
|
||||
val resultChannel = Channel<CountResult>(UNLIMITED)
|
||||
@@ -73,7 +73,7 @@ suspend fun INostrClient.count(
|
||||
|
||||
count(subId = subId, filters = mapOf(relay to listOf(filter)))
|
||||
|
||||
withTimeoutOrNull(timeoutMs) {
|
||||
withTimeoutOrNull(idleTimeoutMs) {
|
||||
resultChannel.receive()
|
||||
}
|
||||
} finally {
|
||||
@@ -91,7 +91,7 @@ suspend fun INostrClient.count(
|
||||
* (one filter per relay) and suspends until all results arrive
|
||||
* or the timeout expires.
|
||||
*
|
||||
* [timeoutMs] is an **idle window measured from the most recent progress**, not a
|
||||
* [idleTimeoutMs] is an **idle window measured from the most recent progress**, not a
|
||||
* wall-clock deadline for the whole batch — the package-wide accessory
|
||||
* convention: each *new* relay's COUNT result restarts it, so a large fan-out
|
||||
* where results keep trickling in is never cut short. A relay re-sending a result
|
||||
@@ -101,12 +101,12 @@ suspend fun INostrClient.count(
|
||||
* discarding the partial map, which is why this returns whatever arrived instead.
|
||||
*
|
||||
* @param filters Map of relay -> filter to count.
|
||||
* @param timeoutMs Idle window between new responses (default 15 s).
|
||||
* @param idleTimeoutMs Idle window between new responses (default 15 s).
|
||||
* @return Map of relay -> [CountResult] for every relay that responded in time.
|
||||
*/
|
||||
suspend fun INostrClient.count(
|
||||
filters: Map<NormalizedRelayUrl, List<Filter>>,
|
||||
timeoutMs: Long = 15_000,
|
||||
idleTimeoutMs: Long = 15_000,
|
||||
): Map<NormalizedRelayUrl, CountResult> {
|
||||
if (filters.isEmpty()) return emptyMap()
|
||||
|
||||
@@ -144,7 +144,7 @@ suspend fun INostrClient.count(
|
||||
// window per relay without needing a wall-clock ceiling.
|
||||
while (results.size < filters.size) {
|
||||
val progressed =
|
||||
withTimeoutOrNull(timeoutMs) {
|
||||
withTimeoutOrNull(idleTimeoutMs) {
|
||||
while (true) {
|
||||
val (relay, result) = resultChannel.receive()
|
||||
// put() returns the previous value: null means this relay
|
||||
@@ -177,20 +177,20 @@ suspend fun INostrClient.count(
|
||||
*
|
||||
* @param relays List of relays to query.
|
||||
* @param filter The filter to count against.
|
||||
* @param timeoutMs Idle window between responses (default 15 s) — see [count].
|
||||
* @param idleTimeoutMs Idle window between responses (default 15 s) — see [count].
|
||||
* @return A merged [CountResult], or `null` if no relay responded.
|
||||
*/
|
||||
suspend fun INostrClient.countMerged(
|
||||
relays: List<NormalizedRelayUrl>,
|
||||
filter: Filter,
|
||||
timeoutMs: Long = 15_000,
|
||||
idleTimeoutMs: Long = 15_000,
|
||||
): CountResult? {
|
||||
if (relays.isEmpty()) return null
|
||||
|
||||
val results =
|
||||
count(
|
||||
filters = relays.associateWith { listOf(filter) },
|
||||
timeoutMs = timeoutMs,
|
||||
idleTimeoutMs = idleTimeoutMs,
|
||||
)
|
||||
|
||||
if (results.isEmpty()) return null
|
||||
|
||||
+17
-17
@@ -31,47 +31,47 @@ import com.vitorpamplona.quartz.nip01Core.relay.normalizer.RelayUrlNormalizer
|
||||
suspend fun INostrClient.fetchAll(
|
||||
relay: String,
|
||||
filter: Filter,
|
||||
timeoutMs: Long = 30_000L,
|
||||
) = fetchAll(newSubId(), mapOf(RelayUrlNormalizer.normalize(relay) to listOf(filter)), timeoutMs)
|
||||
idleTimeoutMs: Long = 30_000L,
|
||||
) = fetchAll(newSubId(), mapOf(RelayUrlNormalizer.normalize(relay) to listOf(filter)), idleTimeoutMs)
|
||||
|
||||
suspend fun INostrClient.fetchAll(
|
||||
relay: String,
|
||||
filters: List<Filter>,
|
||||
timeoutMs: Long = 30_000L,
|
||||
) = fetchAll(newSubId(), mapOf(RelayUrlNormalizer.normalize(relay) to filters), timeoutMs)
|
||||
idleTimeoutMs: Long = 30_000L,
|
||||
) = fetchAll(newSubId(), mapOf(RelayUrlNormalizer.normalize(relay) to filters), idleTimeoutMs)
|
||||
|
||||
suspend fun INostrClient.fetchAll(
|
||||
subscriptionId: String = newSubId(),
|
||||
relay: String,
|
||||
filters: List<Filter>,
|
||||
timeoutMs: Long = 30_000L,
|
||||
) = fetchAll(subscriptionId, mapOf(RelayUrlNormalizer.normalize(relay) to filters), timeoutMs)
|
||||
idleTimeoutMs: Long = 30_000L,
|
||||
) = fetchAll(subscriptionId, mapOf(RelayUrlNormalizer.normalize(relay) to filters), idleTimeoutMs)
|
||||
|
||||
suspend fun INostrClient.fetchAll(
|
||||
relay: NormalizedRelayUrl,
|
||||
filter: Filter,
|
||||
timeoutMs: Long = 30_000L,
|
||||
) = fetchAll(newSubId(), mapOf(relay to listOf(filter)), timeoutMs)
|
||||
idleTimeoutMs: Long = 30_000L,
|
||||
) = fetchAll(newSubId(), mapOf(relay to listOf(filter)), idleTimeoutMs)
|
||||
|
||||
suspend fun INostrClient.fetchAll(
|
||||
relay: NormalizedRelayUrl,
|
||||
filters: List<Filter>,
|
||||
timeoutMs: Long = 30_000L,
|
||||
) = fetchAll(newSubId(), mapOf(relay to filters), timeoutMs)
|
||||
idleTimeoutMs: Long = 30_000L,
|
||||
) = fetchAll(newSubId(), mapOf(relay to filters), idleTimeoutMs)
|
||||
|
||||
suspend fun INostrClient.fetchAll(
|
||||
subscriptionId: String = newSubId(),
|
||||
relay: NormalizedRelayUrl,
|
||||
filters: List<Filter>,
|
||||
timeoutMs: Long = 30_000L,
|
||||
) = fetchAll(subscriptionId, mapOf(relay to filters), timeoutMs)
|
||||
idleTimeoutMs: Long = 30_000L,
|
||||
) = fetchAll(subscriptionId, mapOf(relay to filters), idleTimeoutMs)
|
||||
|
||||
/**
|
||||
* Subscribe [filters], collect every (deduped) event, and return once every
|
||||
* relay reached a terminal state (EOSE, CLOSED, or cannot-connect) or the
|
||||
* line went quiet for [timeoutMs].
|
||||
* line went quiet for [idleTimeoutMs].
|
||||
*
|
||||
* [timeoutMs] is an **idle window, not a hard cap**: every arriving event or
|
||||
* [idleTimeoutMs] is an **idle window, not a hard cap**: every arriving event or
|
||||
* terminal signal resets it, so a slow relay actively streaming a large
|
||||
* backlog is never cropped mid-delivery. The fetch only gives up after a full
|
||||
* window of silence — or at the [maxTotalMs] wall-clock ceiling (default 10x
|
||||
@@ -85,13 +85,13 @@ suspend fun INostrClient.fetchAll(
|
||||
suspend fun INostrClient.fetchAll(
|
||||
subscriptionId: String = newSubId(),
|
||||
filters: Map<NormalizedRelayUrl, List<Filter>>,
|
||||
timeoutMs: Long = 30_000L,
|
||||
maxTotalMs: Long = timeoutMs * 10,
|
||||
idleTimeoutMs: Long = 30_000L,
|
||||
maxTotalMs: Long = idleTimeoutMs * 10,
|
||||
): List<Event> {
|
||||
val seenIds = mutableSetOf<HexKey>()
|
||||
return fetchAllWithHooks(
|
||||
filters = filters,
|
||||
timeoutMs = timeoutMs,
|
||||
idleTimeoutMs = idleTimeoutMs,
|
||||
subscriptionId = subscriptionId,
|
||||
maxTotalMs = maxTotalMs,
|
||||
) { _, event -> seenIds.add(event.id) }
|
||||
|
||||
+6
-6
@@ -71,7 +71,7 @@ import kotlin.coroutines.coroutineContext
|
||||
*
|
||||
* @param relay The relay to query.
|
||||
* @param filters Filters to apply on every page (the `until` field is overwritten per page).
|
||||
* @param timeoutMs Idle window per page — like every accessory timeout, it is measured
|
||||
* @param idleTimeoutMs Idle window per page — like every accessory timeout, it is measured
|
||||
* from the relay's **most recent message**, not from the page's start: every arriving
|
||||
* event resets it, so a slow relay actively streaming a large page is never cropped
|
||||
* mid-delivery. A page only gives up after this much silence without an EOSE.
|
||||
@@ -92,7 +92,7 @@ import kotlin.coroutines.coroutineContext
|
||||
suspend fun INostrClient.fetchAllPages(
|
||||
relay: NormalizedRelayUrl,
|
||||
filters: List<Filter>,
|
||||
timeoutMs: Long = 30_000L,
|
||||
idleTimeoutMs: Long = 30_000L,
|
||||
onNewPage: ((Long) -> Unit)? = null,
|
||||
onEvent: (Event) -> Unit,
|
||||
): Int {
|
||||
@@ -255,9 +255,9 @@ suspend fun INostrClient.fetchAllPages(
|
||||
subscribe(subId, mapOf(relay to activeFilters.map { it.value }), listener)
|
||||
|
||||
// Wait for the page's terminal signal (EOSE / CLOSED / cannot-connect),
|
||||
// giving up only after [timeoutMs] of silence — the wait resets on every
|
||||
// giving up only after [idleTimeoutMs] of silence — the wait resets on every
|
||||
// arriving event, so an actively streaming page is never cut mid-delivery.
|
||||
doneChannel.receiveWithinIdle(clock, timeoutMs)
|
||||
doneChannel.receiveWithinIdle(clock, idleTimeoutMs)
|
||||
|
||||
unsubscribe(subId)
|
||||
doneChannel.close()
|
||||
@@ -309,14 +309,14 @@ suspend fun INostrClient.fetchAllPages(
|
||||
suspend fun INostrClient.fetchAllPages(
|
||||
relay: String,
|
||||
filters: List<Filter>,
|
||||
timeoutMs: Long = 30_000L,
|
||||
idleTimeoutMs: Long = 30_000L,
|
||||
onNewPage: ((Long) -> Unit)? = null,
|
||||
onEvent: (Event) -> Unit,
|
||||
): Int =
|
||||
fetchAllPages(
|
||||
relay = RelayUrlNormalizer.normalize(relay),
|
||||
filters = filters,
|
||||
timeoutMs = timeoutMs,
|
||||
idleTimeoutMs = idleTimeoutMs,
|
||||
onNewPage = onNewPage,
|
||||
onEvent = onEvent,
|
||||
)
|
||||
|
||||
+3
-3
@@ -50,7 +50,7 @@ import kotlinx.coroutines.sync.Semaphore
|
||||
* @param filters per-relay filter lists; the key set is the relays queried, in
|
||||
* iteration order (pass a [LinkedHashMap]/`associateWith` result to control it).
|
||||
* A `search` filter is fetched as a single relevance page — see [fetchAllPages].
|
||||
* @param timeoutMs per-page idle window handed to each relay's [fetchAllPages] —
|
||||
* @param idleTimeoutMs per-page idle window handed to each relay's [fetchAllPages] —
|
||||
* measured from that relay's most recent message (every event resets it), not
|
||||
* from the page's start. As in [fetchAllPages] there is no wall-clock ceiling;
|
||||
* bound a relay's walk with a [Filter.limit], or cancel the caller.
|
||||
@@ -63,7 +63,7 @@ import kotlinx.coroutines.sync.Semaphore
|
||||
*/
|
||||
suspend fun INostrClient.fetchAllPagesFromPool(
|
||||
filters: Map<NormalizedRelayUrl, List<Filter>>,
|
||||
timeoutMs: Long = 30_000L,
|
||||
idleTimeoutMs: Long = 30_000L,
|
||||
maxConcurrentRelays: Int = 8,
|
||||
onNewPage: ((until: Long, relay: NormalizedRelayUrl) -> Unit)? = null,
|
||||
onRelayStart: ((relay: NormalizedRelayUrl) -> Unit)? = null,
|
||||
@@ -83,7 +83,7 @@ suspend fun INostrClient.fetchAllPagesFromPool(
|
||||
fetchAllPages(
|
||||
relay = relay,
|
||||
filters = filtersForRelay,
|
||||
timeoutMs = timeoutMs,
|
||||
idleTimeoutMs = idleTimeoutMs,
|
||||
onNewPage = onNewPage?.let { cb -> { until -> cb(until, relay) } },
|
||||
) { event -> onEvent(event, relay) }
|
||||
onRelayComplete?.invoke(relay, total)
|
||||
|
||||
+11
-11
@@ -41,13 +41,13 @@ import kotlinx.coroutines.withTimeoutOrNull
|
||||
* funnel every arriving event through the suspending [onEvent] hook (verify /
|
||||
* persist / filter — return `true` to keep it in the result), and return the
|
||||
* accepted `(relay, event)` pairs once every relay reached a terminal state
|
||||
* (EOSE, CLOSED, or cannot-connect) or the line went quiet for [timeoutMs].
|
||||
* (EOSE, CLOSED, or cannot-connect) or the line went quiet for [idleTimeoutMs].
|
||||
*
|
||||
* [timeoutMs] is an **idle window, not a hard cap**: the clock only runs while
|
||||
* [idleTimeoutMs] is an **idle window, not a hard cap**: the clock only runs while
|
||||
* the relays are silent, and every arriving event or terminal signal resets
|
||||
* it. A slow relay actively streaming a large backlog is therefore never
|
||||
* cropped mid-delivery — the fetch ends when the work is done or when nothing
|
||||
* has arrived for [timeoutMs] (a stall). The terminal conditions (EOSE /
|
||||
* has arrived for [idleTimeoutMs] (a stall). The terminal conditions (EOSE /
|
||||
* CLOSED / cannot-connect per relay) are what bound the fetch; the timeout's
|
||||
* only job is detecting relays that will never reach one. [maxTotalMs]
|
||||
* (default 10x the idle window) is the wall-clock ceiling that keeps a
|
||||
@@ -62,20 +62,20 @@ import kotlinx.coroutines.withTimeoutOrNull
|
||||
* as a hard failure via [classifyDrainFailure] (connect refused / DNS / TLS /
|
||||
* dead HTTP upgrade — NOT slow relays or 429s) is recorded, so callers can
|
||||
* prune proven-dead relays from future routing instead of paying the full
|
||||
* [timeoutMs] on them again.
|
||||
* [idleTimeoutMs] on them again.
|
||||
* - **[pendingOnAuthRequired]** — a relay that refuses the REQ with an
|
||||
* `auth-required:` CLOSED is kept pending rather than treated as terminal:
|
||||
* the caller's NIP-42 responder answers the challenge and the client re-fires
|
||||
* this same subscription, so the post-auth events are collected instead of
|
||||
* returning empty. If auth never satisfies it, the relay simply falls through
|
||||
* to the [timeoutMs].
|
||||
* to the [idleTimeoutMs].
|
||||
* - **[onTimeout]** — diagnostic hook fired when the idle window elapsed with
|
||||
* relays still pending: receives the stalled set, the terminal reasons seen so
|
||||
* far (`"eose"` / `"closed:<msg>"` / `"cannot:<msg>"`), and what was collected.
|
||||
*/
|
||||
suspend fun INostrClient.fetchAllWithHooks(
|
||||
filters: Map<NormalizedRelayUrl, List<Filter>>,
|
||||
timeoutMs: Long = 8_000L,
|
||||
idleTimeoutMs: Long = 8_000L,
|
||||
subscriptionId: String = newSubId(),
|
||||
pendingOnAuthRequired: Boolean = false,
|
||||
deadOut: MutableMap<NormalizedRelayUrl, DrainFailure>? = null,
|
||||
@@ -87,10 +87,10 @@ suspend fun INostrClient.fetchAllWithHooks(
|
||||
* restores an upper bound while staying far above the idle window, so a
|
||||
* legitimately streaming relay still finishes its backlog. Pass
|
||||
* [Long.MAX_VALUE] for a deliberately uncapped drain; a non-positive value
|
||||
* also uncaps (absorbing a `timeoutMs * 10` overflow from an
|
||||
* also uncaps (absorbing an `idleTimeoutMs * 10` overflow from an
|
||||
* effectively-infinite idle window).
|
||||
*/
|
||||
maxTotalMs: Long = timeoutMs * 10,
|
||||
maxTotalMs: Long = idleTimeoutMs * 10,
|
||||
onEvent: suspend (relay: NormalizedRelayUrl, event: Event) -> Boolean,
|
||||
): List<Pair<NormalizedRelayUrl, Event>> {
|
||||
if (filters.isEmpty()) return emptyList()
|
||||
@@ -190,7 +190,7 @@ suspend fun INostrClient.fetchAllWithHooks(
|
||||
// Slow path: both dry — arm one idle wait for the next signal.
|
||||
if (pending == null) {
|
||||
val progressed =
|
||||
withTimeoutOrNull(timeoutMs) {
|
||||
withTimeoutOrNull(idleTimeoutMs) {
|
||||
select<Unit> {
|
||||
eventChannel.onReceive { pending = it }
|
||||
doneChannel.onReceive { (relay, reason) ->
|
||||
@@ -256,7 +256,7 @@ suspend fun INostrClient.fetchAllWithHooks(
|
||||
*/
|
||||
suspend fun INostrClient.fetchAllPagesFromPoolWithHooks(
|
||||
filters: Map<NormalizedRelayUrl, List<Filter>>,
|
||||
timeoutMs: Long = 30_000L,
|
||||
idleTimeoutMs: Long = 30_000L,
|
||||
maxConcurrentRelays: Int = 8,
|
||||
onEvent: suspend (relay: NormalizedRelayUrl, event: Event) -> Boolean,
|
||||
): List<Pair<NormalizedRelayUrl, Event>> {
|
||||
@@ -287,7 +287,7 @@ suspend fun INostrClient.fetchAllPagesFromPoolWithHooks(
|
||||
try {
|
||||
fetchAllPagesFromPool(
|
||||
filters = filters,
|
||||
timeoutMs = timeoutMs,
|
||||
idleTimeoutMs = idleTimeoutMs,
|
||||
maxConcurrentRelays = maxConcurrentRelays,
|
||||
) { event, relay -> eventChannel.trySend(relay to event) }
|
||||
} finally {
|
||||
|
||||
+3
-3
@@ -69,7 +69,7 @@ suspend fun INostrClient.fetchFirst(
|
||||
* every relay reached a terminal state — EOSE, CLOSED, or cannot-connect — with
|
||||
* nothing matching, or the line went quiet).
|
||||
*
|
||||
* [timeoutMs] is an **idle window measured from the most recent progress**, not a
|
||||
* [idleTimeoutMs] is an **idle window measured from the most recent progress**, not a
|
||||
* wall-clock deadline — the package-wide accessory convention. Progress means a
|
||||
* signal that actually advances the fetch: an event, or the first terminal state
|
||||
* from a relay still being waited on. Repeat chatter from a relay already
|
||||
@@ -86,7 +86,7 @@ suspend fun INostrClient.fetchFirst(
|
||||
suspend fun INostrClient.fetchFirst(
|
||||
subscriptionId: String = newSubId(),
|
||||
filters: Map<NormalizedRelayUrl, List<Filter>>,
|
||||
timeoutMs: Long = 30_000L,
|
||||
idleTimeoutMs: Long = 30_000L,
|
||||
): Event? {
|
||||
val eventChannel = Channel<Event>(UNLIMITED)
|
||||
val doneChannel = Channel<NormalizedRelayUrl>(UNLIMITED)
|
||||
@@ -137,7 +137,7 @@ suspend fun INostrClient.fetchFirst(
|
||||
// advance escapes to the outer loop and earns a fresh window.
|
||||
while (remaining.isNotEmpty()) {
|
||||
val progressed =
|
||||
withTimeoutOrNull(timeoutMs) {
|
||||
withTimeoutOrNull(idleTimeoutMs) {
|
||||
while (true) {
|
||||
val advanced =
|
||||
select<Boolean> {
|
||||
|
||||
+14
-9
@@ -14,11 +14,16 @@ Import as `com.vitorpamplona.quartz.nip01Core.relay.client.accessories.<name>` (
|
||||
|
||||
## Timeout convention
|
||||
|
||||
Every `timeoutMs` / `idleTimeoutMs` in this package is an **idle window measured
|
||||
from the relay's most recent progress**, never a wall-clock deadline: real progress
|
||||
resets it, so an actively streaming relay is never cut off mid-delivery — the
|
||||
operation only gives up after a full window of silence. The shared primitives are in
|
||||
`IdleWatchdog.kt` (`IdleClock` + `receiveWithinIdle`); use them in a new accessory.
|
||||
Every wait in this package is an **idle window measured from the relay's most recent
|
||||
progress**, never a wall-clock deadline: real progress resets it, so an actively
|
||||
streaming relay is never cut off mid-delivery — the operation only gives up after a
|
||||
full window of silence. The shared primitives are in `IdleWatchdog.kt` (`IdleClock` +
|
||||
`receiveWithinIdle`); use them in a new accessory.
|
||||
|
||||
**The parameter is named `idleTimeoutMs`, never `timeoutMs`** — the name is the
|
||||
contract, so a caller can't mistake it for a deadline. The sole exception is
|
||||
`publishAndConfirm`'s `timeoutInSeconds`, which genuinely *is* a fixed window (see
|
||||
below); the differing name is the tell.
|
||||
|
||||
**Progress, not merely traffic.** A message that tells us nothing new — a relay
|
||||
re-CLOSEing after we already recorded it as done, a duplicate COUNT — must not
|
||||
@@ -53,9 +58,9 @@ window to collect the `OK`s — a bounded confirmation round-trip, not a stream.
|
||||
|
||||
| Function | File | Use when |
|
||||
| --- | --- | --- |
|
||||
| `fetchAll(relay, filter, timeoutMs)` | `NostrClientFetchAllExt` | Get every event matching a filter in one REQ, deduped by id, until EOSE or a full idle window of silence. **No verify, no store** — just the events. |
|
||||
| `fetchFirst(relay, filter, timeoutMs)` | `NostrClientFetchFirstExt` | Get the first matching event and stop (returns `null` on none/timeout). |
|
||||
| `fetchAllPages(relay, filters, timeoutMs)` | `NostrClientFetchAllPagesExt` | Fully retrieve a result set larger than the relay's per-REQ cap (strfry `limit`, ~500) by walking a `created_at` cursor. Bound it with the filter's `limit`. |
|
||||
| `fetchAll(relay, filter, idleTimeoutMs)` | `NostrClientFetchAllExt` | Get every event matching a filter in one REQ, deduped by id, until EOSE or a full idle window of silence. **No verify, no store** — just the events. |
|
||||
| `fetchFirst(relay, filter, idleTimeoutMs)` | `NostrClientFetchFirstExt` | Get the first matching event and stop (returns `null` on none/timeout). |
|
||||
| `fetchAllPages(relay, filters, idleTimeoutMs)` | `NostrClientFetchAllPagesExt` | Fully retrieve a result set larger than the relay's per-REQ cap (strfry `limit`, ~500) by walking a `created_at` cursor. Bound it with the filter's `limit`. |
|
||||
| `fetchAllPagesFromPool(filters, ...)` | `NostrClientFetchAllPagesPoolExt` | Same paging, across several relays at once. No cross-relay dedup — the `WithHooks` variant below dedups. |
|
||||
| `fetchAllWithHooks(filters, ...)` | `NostrClientFetchAllWithHooksExt` | `fetchAll` with a suspending per-`(relay, event)` accept hook (verify+store as events arrive), per-relay terminal-reason tracking, optional dead-relay collection (`deadOut` + `classifyDrainFailure`), keep-pending-on-`auth-required` CLOSED (NIP-42 re-fire), and a timeout diagnostic hook. |
|
||||
| `fetchAllPagesFromPoolWithHooks(filters, ...)` | `NostrClientFetchAllWithHooksExt` | `fetchAllPagesFromPool` with the same suspending accept hook, run single-threaded in one consumer; deduped across relays by `SeenIds` before the hook. |
|
||||
@@ -79,7 +84,7 @@ window to collect the `OK`s — a bounded confirmation round-trip, not a stream.
|
||||
|
||||
| Function | File | Use when |
|
||||
| --- | --- | --- |
|
||||
| `count(relay, filter, timeoutMs)` | `NostrClientCountExt` | NIP-45 `COUNT` against one relay (`null` on timeout / no support). |
|
||||
| `count(relay, filter, idleTimeoutMs)` | `NostrClientCountExt` | NIP-45 `COUNT` against one relay (`null` on timeout / no support). |
|
||||
| `countMerged(relays, filter, ...)` | `NostrClientCountExt` | Merged count across relays. |
|
||||
|
||||
## Negentropy (NIP-77)
|
||||
|
||||
+5
-5
@@ -89,7 +89,7 @@ class FetchAllIdleTimeoutTest {
|
||||
val collected =
|
||||
client.fetchAllWithHooks(
|
||||
filters = mapOf(relay to listOf(Filter(kinds = listOf(1)))),
|
||||
timeoutMs = 300,
|
||||
idleTimeoutMs = 300,
|
||||
) { _, _ -> true }
|
||||
feeder.join()
|
||||
assertEquals(10, collected.size, "an actively streaming relay must never be cropped")
|
||||
@@ -111,7 +111,7 @@ class FetchAllIdleTimeoutTest {
|
||||
val collected =
|
||||
client.fetchAllWithHooks(
|
||||
filters = mapOf(relay to listOf(Filter(kinds = listOf(1)))),
|
||||
timeoutMs = 300,
|
||||
idleTimeoutMs = 300,
|
||||
onTimeout = { stalled, _, _ -> stalledRelays = stalled },
|
||||
) { _, _ -> true }
|
||||
assertEquals(2, collected.size, "events before the stall are kept")
|
||||
@@ -136,7 +136,7 @@ class FetchAllIdleTimeoutTest {
|
||||
val events =
|
||||
client.fetchAll(
|
||||
filters = mapOf(relay to listOf(Filter(kinds = listOf(1)))),
|
||||
timeoutMs = 300,
|
||||
idleTimeoutMs = 300,
|
||||
)
|
||||
feeder.join()
|
||||
assertEquals(10, events.size, "fetchAll shares the idle-window semantics")
|
||||
@@ -162,7 +162,7 @@ class FetchAllIdleTimeoutTest {
|
||||
val collected =
|
||||
client.fetchAllWithHooks(
|
||||
filters = mapOf(relay to listOf(Filter(kinds = listOf(1)))),
|
||||
timeoutMs = 300,
|
||||
idleTimeoutMs = 300,
|
||||
maxTotalMs = 1_000,
|
||||
onTimeout = { stalled, _, _ -> stalledRelays = stalled },
|
||||
) { _, _ -> true }
|
||||
@@ -186,7 +186,7 @@ class FetchAllIdleTimeoutTest {
|
||||
val collected =
|
||||
client.fetchAllWithHooks(
|
||||
filters = mapOf(relay to listOf(Filter(kinds = listOf(1)))),
|
||||
timeoutMs = 300,
|
||||
idleTimeoutMs = 300,
|
||||
) { _, _ -> true }
|
||||
assertEquals(1, collected.size)
|
||||
assertTrue(currentTime - start < 300, "a terminal EOSE must not wait out the window")
|
||||
|
||||
+10
-9
@@ -38,10 +38,11 @@ import kotlin.test.assertEquals
|
||||
import kotlin.test.assertNull
|
||||
|
||||
/**
|
||||
* Pins [fetchFirst]'s timeout to the package-wide idle-window convention:
|
||||
* [timeoutMs] is silence measured from the most recent relay signal, not an
|
||||
* absolute deadline across the whole multi-relay wait — with [maxTotalMs] as
|
||||
* the wall-clock ceiling.
|
||||
* Pins [fetchFirst]'s timeout to the package-wide convention: `idleTimeoutMs` is
|
||||
* silence measured from the most recent *progress*, not an absolute deadline
|
||||
* across the whole multi-relay wait. Repeat chatter from a relay already
|
||||
* accounted for is not progress, which is what makes the call self-bounding
|
||||
* without a ceiling parameter — a hard bound is the caller's `withTimeoutOrNull`.
|
||||
*/
|
||||
@OptIn(ExperimentalCoroutinesApi::class)
|
||||
class FetchFirstIdleTimeoutTest {
|
||||
@@ -97,7 +98,7 @@ class FetchFirstIdleTimeoutTest {
|
||||
val result =
|
||||
client.fetchFirst(
|
||||
filters = filters(relayA, relayB, relayC, relayD),
|
||||
timeoutMs = 300,
|
||||
idleTimeoutMs = 300,
|
||||
)
|
||||
assertEquals(event(1).id, result?.id, "progress must restart the window; the slow relay's event still lands")
|
||||
}
|
||||
@@ -110,7 +111,7 @@ class FetchFirstIdleTimeoutTest {
|
||||
val result =
|
||||
client.fetchFirst(
|
||||
filters = filters(relayA),
|
||||
timeoutMs = 300,
|
||||
idleTimeoutMs = 300,
|
||||
)
|
||||
assertNull(result)
|
||||
assertEquals(300L, currentTime - start, "a silent relay costs exactly one idle window")
|
||||
@@ -136,7 +137,7 @@ class FetchFirstIdleTimeoutTest {
|
||||
val result =
|
||||
client.fetchFirst(
|
||||
filters = filters(relayA, relayB),
|
||||
timeoutMs = 300,
|
||||
idleTimeoutMs = 300,
|
||||
)
|
||||
chatter.cancel()
|
||||
assertNull(result)
|
||||
@@ -160,7 +161,7 @@ class FetchFirstIdleTimeoutTest {
|
||||
val result =
|
||||
client.fetchFirst(
|
||||
filters = filters(relayA),
|
||||
timeoutMs = 300,
|
||||
idleTimeoutMs = 300,
|
||||
)
|
||||
assertEquals(event(7).id, result?.id, "an event racing the final EOSE must not be dropped")
|
||||
}
|
||||
@@ -180,7 +181,7 @@ class FetchFirstIdleTimeoutTest {
|
||||
// is the wall-clock bound, and costs nothing because a timed-out
|
||||
// fetchFirst yields null either way.
|
||||
val start = currentTime
|
||||
val result = withTimeoutOrNull(120) { client.fetchFirst(filters = filters(relayA, relayB), timeoutMs = 10_000) }
|
||||
val result = withTimeoutOrNull(120) { client.fetchFirst(filters = filters(relayA, relayB), idleTimeoutMs = 10_000) }
|
||||
chatter.cancel()
|
||||
assertNull(result)
|
||||
assertEquals(120L, currentTime - start, "the caller's timeout bounds the call")
|
||||
|
||||
+6
-7
@@ -40,7 +40,7 @@ import kotlin.time.TimeSource
|
||||
|
||||
/**
|
||||
* Pins [fetchAllPages]'s timeout to the package-wide idle-window convention:
|
||||
* `timeoutMs` is silence measured from the relay's MOST RECENT message, not a
|
||||
* `idleTimeoutMs` is silence measured from the relay's MOST RECENT message, not a
|
||||
* wall-clock deadline for the page. A relay that keeps streaming — however
|
||||
* slowly — must never have a page cropped mid-delivery.
|
||||
*
|
||||
@@ -106,7 +106,7 @@ class NostrClientFetchAllPagesIdleTimeoutTest {
|
||||
client.fetchAllPages(
|
||||
relay = relay,
|
||||
filters = listOf(Filter(kinds = listOf(1), limit = 6)),
|
||||
timeoutMs = 500,
|
||||
idleTimeoutMs = 500,
|
||||
onNewPage = { pages++ },
|
||||
) { got.add(it) }
|
||||
feeder.join()
|
||||
@@ -140,7 +140,7 @@ class NostrClientFetchAllPagesIdleTimeoutTest {
|
||||
client.fetchAllPages(
|
||||
relay = relay,
|
||||
filters = listOf(Filter(kinds = listOf(1), limit = 3)),
|
||||
timeoutMs = 300,
|
||||
idleTimeoutMs = 300,
|
||||
) { got.add(it) }
|
||||
feeder.join()
|
||||
val elapsedMs = start.elapsedNow().inWholeMilliseconds
|
||||
@@ -151,9 +151,8 @@ class NostrClientFetchAllPagesIdleTimeoutTest {
|
||||
}
|
||||
|
||||
/**
|
||||
* Documents why [fetchAllPages] has no wall-clock ceiling, unlike the
|
||||
* single-wait accessories (`fetchAll`/`fetchFirst`/`count`, whose `maxTotalMs`
|
||||
* really does end the call).
|
||||
* Documents why [fetchAllPages] has no wall-clock ceiling — nor should any
|
||||
* accessory: a hard bound composes at the call site as `withTimeoutOrNull`.
|
||||
*
|
||||
* A per-page ceiling cannot bound this walk: when a page ends, the loop advances
|
||||
* the cursor and fires the NEXT `REQ`, so an endless trickle against an unbounded
|
||||
@@ -181,7 +180,7 @@ class NostrClientFetchAllPagesIdleTimeoutTest {
|
||||
client.fetchAllPages(
|
||||
relay = relay,
|
||||
filters = listOf(Filter(kinds = listOf(1))), // unbounded: no limit
|
||||
timeoutMs = 200,
|
||||
idleTimeoutMs = 200,
|
||||
) { }
|
||||
}
|
||||
feeder.cancel()
|
||||
|
||||
+3
-3
@@ -141,7 +141,7 @@ class BulkDownloadBenchmark {
|
||||
client.fetchAllPages(
|
||||
relay = relayUrl,
|
||||
filters = listOf(Filter(kinds = listOf(1), since = lo, until = hi)),
|
||||
timeoutMs = LOCAL_PAGE_TIMEOUT_MS,
|
||||
idleTimeoutMs = LOCAL_PAGE_TIMEOUT_MS,
|
||||
) { event ->
|
||||
count.incrementAndGet()
|
||||
bytes.addAndGet(event.content.length.toLong())
|
||||
@@ -376,7 +376,7 @@ class BulkDownloadBenchmark {
|
||||
client.fetchAllPages(
|
||||
relay = relay,
|
||||
filters = listOf(Filter(kinds = listOf(PROD_KIND), limit = PROD_MAX_EVENTS)),
|
||||
timeoutMs = 30_000L,
|
||||
idleTimeoutMs = 30_000L,
|
||||
onNewPage = { pages++ },
|
||||
) { event ->
|
||||
count.incrementAndGet()
|
||||
@@ -672,7 +672,7 @@ class BulkDownloadBenchmark {
|
||||
client.fetchAllPages(
|
||||
relay = relay,
|
||||
filters = listOf(Filter(kinds = listOf(PROD_KIND), limit = PROD_MAX_EVENTS)),
|
||||
timeoutMs = 30_000L,
|
||||
idleTimeoutMs = 30_000L,
|
||||
onNewPage = { pages++ },
|
||||
) { event ->
|
||||
count.incrementAndGet()
|
||||
|
||||
+1
-1
@@ -280,7 +280,7 @@ class ByIdFetchBenchmark {
|
||||
client.fetchAllPages(
|
||||
relay = relay,
|
||||
filters = listOf(Filter(kinds = listOf(KIND))),
|
||||
timeoutMs = ENUM_PAGE_TIMEOUT_MS,
|
||||
idleTimeoutMs = ENUM_PAGE_TIMEOUT_MS,
|
||||
) { event -> ids.add(event.id) }
|
||||
println(" enumerated ${ids.size} ids in %.1fs (paged, untimed baseline for the matrix)".format((System.nanoTime() - t) / 1e9))
|
||||
ids.toList()
|
||||
|
||||
Reference in New Issue
Block a user