mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-05 19:28:25 +00:00
refactor: drop fetchAllPages' maxPageMs ceiling
The per-page wall-clock ceiling added alongside the idle-window change did not do what it claimed. fetchAllPages waits inside a while(true): when a page's wait ends, the loop advances the until cursor and issues the next REQ rather than returning, so a ceiling bounds one page, never the call. Measured against a relay trickling events forever with an unbounded filter, maxPageMs=400 produced 8 REQs and no return -- the walk ran until the caller cancelled, exactly as it would with no ceiling at all. It also made truncation unsafe. Cutting a page mid-stream advances until to the oldest event received so far, which only preserves the set if the relay streams strictly newest-first -- NIP-01 recommends that but does not require it -- so a ceiling firing on an out-of-order relay can skip the not-yet-sent events above that cursor. That is the same class of silent gap the idle window was introduced to close. What actually bounds a paged download is the filter's limit (already the documented way) or cancelling the caller, which the ensureActive() at the top of each page honors. Removed from fetchAllPages, fetchAllPagesFromPool and fetchAllPagesFromPoolWithHooks; a regression test pins the real contract so nobody re-adds a ceiling believing it caps the walk. maxTotalMs stays on fetchAll/fetchAllWithHooks, fetchFirst and multi-relay count: those wait in a single loop and return, so there the ceiling genuinely ends the call (each is covered by a test asserting it does). Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01KZe1Pq2ejoHehdufhPDRgb
This commit is contained in:
+13
-15
@@ -30,7 +30,6 @@ import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.RelayUrlNormalizer
|
||||
import kotlinx.coroutines.channels.Channel
|
||||
import kotlinx.coroutines.ensureActive
|
||||
import kotlinx.coroutines.withTimeoutOrNull
|
||||
import kotlin.coroutines.coroutineContext
|
||||
|
||||
/**
|
||||
@@ -76,9 +75,17 @@ import kotlin.coroutines.coroutineContext
|
||||
* 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.
|
||||
* @param maxPageMs Hard wall-clock ceiling per page (default 10x the idle window) — the
|
||||
* idle window alone is unbounded against a misbehaving relay that trickles events
|
||||
* forever without ever sending EOSE. Pass [Long.MAX_VALUE] for an uncapped page.
|
||||
*
|
||||
* Deliberately no wall-clock ceiling here, unlike [fetchAll]'s `maxTotalMs`. A ceiling
|
||||
* would bound one *page*, not this call: the loop below reacts to a page ending by
|
||||
* advancing the cursor and issuing the next REQ, so a relay trickling events forever
|
||||
* against an unbounded filter would just be re-paged forever — measurably so (see
|
||||
* NostrClientFetchAllPagesIdleTimeoutTest). Worse, cutting a page mid-stream advances
|
||||
* `until` to the oldest event received *so far*, which only preserves the set if the
|
||||
* relay streams strictly newest-first (NIP-01 recommends but does not require it) —
|
||||
* otherwise the not-yet-sent events above that cursor are skipped. What actually
|
||||
* bounds this walk is a [Filter.limit] (the documented way to cap a download) or
|
||||
* cancelling the caller, which the [ensureActive] at the top of each page honors.
|
||||
* @param onEvent Called once for every distinct event delivered, in page order.
|
||||
* @return Total number of distinct events delivered across all pages.
|
||||
*/
|
||||
@@ -86,7 +93,6 @@ suspend fun INostrClient.fetchAllPages(
|
||||
relay: NormalizedRelayUrl,
|
||||
filters: List<Filter>,
|
||||
timeoutMs: Long = 30_000L,
|
||||
maxPageMs: Long = timeoutMs * 10,
|
||||
onNewPage: ((Long) -> Unit)? = null,
|
||||
onEvent: (Event) -> Unit,
|
||||
): Int {
|
||||
@@ -237,14 +243,8 @@ suspend fun INostrClient.fetchAllPages(
|
||||
|
||||
// Wait for the page's terminal signal (EOSE / CLOSED / cannot-connect),
|
||||
// giving up only after [timeoutMs] of silence — the wait resets on every
|
||||
// event — or at the [maxPageMs] wall-clock ceiling. A non-positive
|
||||
// ceiling means uncapped, mirroring the idle window's `<= 0` = disabled
|
||||
// convention (it also absorbs a `timeoutMs * 10` overflow from a caller
|
||||
// passing an effectively-infinite idle window).
|
||||
val ceiling = if (maxPageMs <= 0) Long.MAX_VALUE else maxPageMs
|
||||
withTimeoutOrNull(ceiling) {
|
||||
doneChannel.receiveWithinIdle(clock, timeoutMs)
|
||||
}
|
||||
// arriving event, so an actively streaming page is never cut mid-delivery.
|
||||
doneChannel.receiveWithinIdle(clock, timeoutMs)
|
||||
|
||||
unsubscribe(subId)
|
||||
doneChannel.close()
|
||||
@@ -297,7 +297,6 @@ suspend fun INostrClient.fetchAllPages(
|
||||
relay: String,
|
||||
filters: List<Filter>,
|
||||
timeoutMs: Long = 30_000L,
|
||||
maxPageMs: Long = timeoutMs * 10,
|
||||
onNewPage: ((Long) -> Unit)? = null,
|
||||
onEvent: (Event) -> Unit,
|
||||
): Int =
|
||||
@@ -305,7 +304,6 @@ suspend fun INostrClient.fetchAllPages(
|
||||
relay = RelayUrlNormalizer.normalize(relay),
|
||||
filters = filters,
|
||||
timeoutMs = timeoutMs,
|
||||
maxPageMs = maxPageMs,
|
||||
onNewPage = onNewPage,
|
||||
onEvent = onEvent,
|
||||
)
|
||||
|
||||
+2
-4
@@ -52,8 +52,8 @@ import kotlinx.coroutines.sync.Semaphore
|
||||
* A `search` filter is fetched as a single relevance page — see [fetchAllPages].
|
||||
* @param timeoutMs 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.
|
||||
* @param maxPageMs per-page wall-clock ceiling handed to each relay's [fetchAllPages].
|
||||
* 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.
|
||||
* @param maxConcurrentRelays upper bound on relays paginating at once (≥ 1).
|
||||
* @param onNewPage optional `(until, relay)` tick before each non-first page.
|
||||
* @param onRelayStart optional hook fired as each relay's download begins.
|
||||
@@ -64,7 +64,6 @@ import kotlinx.coroutines.sync.Semaphore
|
||||
suspend fun INostrClient.fetchAllPagesFromPool(
|
||||
filters: Map<NormalizedRelayUrl, List<Filter>>,
|
||||
timeoutMs: Long = 30_000L,
|
||||
maxPageMs: Long = timeoutMs * 10,
|
||||
maxConcurrentRelays: Int = 8,
|
||||
onNewPage: ((until: Long, relay: NormalizedRelayUrl) -> Unit)? = null,
|
||||
onRelayStart: ((relay: NormalizedRelayUrl) -> Unit)? = null,
|
||||
@@ -85,7 +84,6 @@ suspend fun INostrClient.fetchAllPagesFromPool(
|
||||
relay = relay,
|
||||
filters = filtersForRelay,
|
||||
timeoutMs = timeoutMs,
|
||||
maxPageMs = maxPageMs,
|
||||
onNewPage = onNewPage?.let { cb -> { until -> cb(until, relay) } },
|
||||
) { event -> onEvent(event, relay) }
|
||||
onRelayComplete?.invoke(relay, total)
|
||||
|
||||
-2
@@ -257,7 +257,6 @@ suspend fun INostrClient.fetchAllWithHooks(
|
||||
suspend fun INostrClient.fetchAllPagesFromPoolWithHooks(
|
||||
filters: Map<NormalizedRelayUrl, List<Filter>>,
|
||||
timeoutMs: Long = 30_000L,
|
||||
maxPageMs: Long = timeoutMs * 10,
|
||||
maxConcurrentRelays: Int = 8,
|
||||
onEvent: suspend (relay: NormalizedRelayUrl, event: Event) -> Boolean,
|
||||
): List<Pair<NormalizedRelayUrl, Event>> {
|
||||
@@ -289,7 +288,6 @@ suspend fun INostrClient.fetchAllPagesFromPoolWithHooks(
|
||||
fetchAllPagesFromPool(
|
||||
filters = filters,
|
||||
timeoutMs = timeoutMs,
|
||||
maxPageMs = maxPageMs,
|
||||
maxConcurrentRelays = maxConcurrentRelays,
|
||||
) { event, relay -> eventChannel.trySend(relay to event) }
|
||||
} finally {
|
||||
|
||||
+18
-5
@@ -19,11 +19,24 @@ from the relay's most recent message**, never a wall-clock deadline: each arrivi
|
||||
event / result / terminal signal 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 when writing a new accessory. Because an idle window alone is unbounded
|
||||
against a relay that trickles messages forever, the loops that could run away also
|
||||
take a hard wall-clock ceiling (`maxTotalMs` / `maxPageMs`, default 10x the idle
|
||||
window) as a backstop; a non-positive ceiling means uncapped. The one deliberate
|
||||
exception is the write side: `publishAndConfirm`'s `timeoutInSeconds` is a fixed
|
||||
use them when writing a new accessory.
|
||||
|
||||
An idle window alone never expires against a relay that trickles messages forever,
|
||||
so the accessories that wait in **one** loop and then return — `fetchAll` /
|
||||
`fetchAllWithHooks`, `fetchFirst`, multi-relay `count` — also take a wall-clock
|
||||
ceiling (`maxTotalMs`, default 10x the idle window; non-positive means uncapped).
|
||||
There the ceiling genuinely ends the call.
|
||||
|
||||
`fetchAllPages` deliberately has **no** ceiling. A per-page ceiling would bound a
|
||||
page, not the call: the paging loop reacts to a page ending by advancing the cursor
|
||||
and issuing the next `REQ`, so an endless trickle just gets re-paged forever (a
|
||||
ceiling of 400 ms against one measured 8 `REQ`s and no return). It also makes
|
||||
truncation unsafe — cutting a page mid-stream advances `until` to the oldest event
|
||||
received *so far*, which only preserves the set if the relay streams strictly
|
||||
newest-first, which NIP-01 recommends but does not require. Bound a paged download
|
||||
with the filter's `limit`, or by cancelling the caller.
|
||||
|
||||
The write side is its own case: `publishAndConfirm`'s `timeoutInSeconds` is a fixed
|
||||
window to collect the `OK`s — a bounded confirmation round-trip, not a stream.
|
||||
|
||||
## One-shot reads (subscribe → collect → return)
|
||||
|
||||
+56
-10
@@ -31,8 +31,10 @@ import com.vitorpamplona.quartz.nip01Core.relay.normalizer.RelayUrlNormalizer
|
||||
import kotlinx.coroutines.delay
|
||||
import kotlinx.coroutines.launch
|
||||
import kotlinx.coroutines.runBlocking
|
||||
import kotlinx.coroutines.withTimeoutOrNull
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertEquals
|
||||
import kotlin.test.assertNull
|
||||
import kotlin.test.assertTrue
|
||||
import kotlin.time.TimeSource
|
||||
|
||||
@@ -49,7 +51,10 @@ import kotlin.time.TimeSource
|
||||
class NostrClientFetchAllPagesIdleTimeoutTest {
|
||||
/** Captures the subscription listener so the test can play a relay. */
|
||||
private class ScriptedClient : INostrClient by EmptyNostrClient() {
|
||||
@Volatile
|
||||
var listener: SubscriptionListener? = null
|
||||
|
||||
@Volatile
|
||||
var subscribeCount = 0
|
||||
|
||||
override fun subscribe(
|
||||
@@ -64,16 +69,18 @@ class NostrClientFetchAllPagesIdleTimeoutTest {
|
||||
|
||||
private val relay = RelayUrlNormalizer.normalize("wss://slow.example.com")
|
||||
|
||||
private fun event(i: Int) =
|
||||
Event(
|
||||
id = i.toString(16).padStart(64, '0'),
|
||||
pubKey = "f".repeat(64),
|
||||
createdAt = i.toLong(),
|
||||
kind = 1,
|
||||
tags = emptyArray(),
|
||||
content = "e$i",
|
||||
sig = "0".repeat(128),
|
||||
)
|
||||
private fun event(
|
||||
i: Int,
|
||||
createdAt: Long = i.toLong(),
|
||||
) = Event(
|
||||
id = i.toString(16).padStart(64, '0'),
|
||||
pubKey = "f".repeat(64),
|
||||
createdAt = createdAt,
|
||||
kind = 1,
|
||||
tags = emptyArray(),
|
||||
content = "e$i",
|
||||
sig = "0".repeat(128),
|
||||
)
|
||||
|
||||
@Test
|
||||
fun streamingPageOutlivesTheIdleWindow() =
|
||||
@@ -142,4 +149,43 @@ class NostrClientFetchAllPagesIdleTimeoutTest {
|
||||
assertTrue(elapsedMs >= 300, "must wait out at least one idle window, took ${elapsedMs}ms")
|
||||
assertTrue(elapsedMs < 5_000, "a stalled page must end promptly after the idle window, took ${elapsedMs}ms")
|
||||
}
|
||||
|
||||
/**
|
||||
* Documents why [fetchAllPages] has no wall-clock ceiling, unlike the
|
||||
* single-wait accessories (`fetchAll`/`fetchFirst`/`count`, whose `maxTotalMs`
|
||||
* really does end the call).
|
||||
*
|
||||
* 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
|
||||
* filter is merely re-paged. This pins that reality — the walk runs until the
|
||||
* caller cancels — so nobody re-adds a `maxPageMs` believing it caps anything.
|
||||
* What actually bounds a download is the filter's `limit`.
|
||||
*/
|
||||
@Test
|
||||
fun anEndlessTrickleIsBoundedByCancellationNotByAWallClock() =
|
||||
runBlocking {
|
||||
val client = ScriptedClient()
|
||||
var ts = 10_000_000L
|
||||
var i = 1
|
||||
val feeder =
|
||||
launch {
|
||||
// Trickles forever, strictly decreasing created_at, never EOSE.
|
||||
while (true) {
|
||||
delay(40)
|
||||
client.listener?.onEvent(event(i++, ts--), false, relay, null)
|
||||
}
|
||||
}
|
||||
|
||||
val returned =
|
||||
withTimeoutOrNull(1_500) {
|
||||
client.fetchAllPages(
|
||||
relay = relay,
|
||||
filters = listOf(Filter(kinds = listOf(1))), // unbounded: no limit
|
||||
timeoutMs = 200,
|
||||
) { }
|
||||
}
|
||||
feeder.cancel()
|
||||
|
||||
assertNull(returned, "an unbounded filter against an endless trickle ends only by cancellation")
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user